Skip to content

Commit d1ad1cb

Browse files
Copilotnnhy
andcommitted
修复Pop/Ack/ChangeInvisibleTime操作缺少queueId参数及添加MessageExt便利方法
Co-authored-by: nnhy <506367+nnhy@users.noreply.github.com>
1 parent 6bfbfe0 commit d1ad1cb

3 files changed

Lines changed: 154 additions & 11 deletions

File tree

NewLife.RocketMQ/Consumer.cs

Lines changed: 40 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1474,12 +1474,13 @@ public Boolean SendMessageBack(MessageExt msg, Int32 delayLevel = 0, Int32 maxRe
14741474
#region Pop消费模式
14751475
/// <summary>Pop方式拉取消息。5.0新增的轻量消费模式,无需客户端Rebalance</summary>
14761476
/// <param name="brokerName">Broker名称</param>
1477+
/// <param name="queueId">队列编号。-1表示由Broker自动分配</param>
14771478
/// <param name="maxNums">最大拉取数</param>
14781479
/// <param name="invisibleTime">不可见时间(毫秒),消息被拉取后在此时间内不会被其他消费者看到</param>
14791480
/// <param name="pollTime">长轮询等待时间(毫秒)</param>
14801481
/// <param name="cancellationToken">取消通知</param>
14811482
/// <returns>拉取结果</returns>
1482-
public async Task<PullResult> PopMessageAsync(String brokerName, Int32 maxNums = 32, Int64 invisibleTime = 60_000, Int32 pollTime = 15_000, CancellationToken cancellationToken = default)
1483+
public async Task<PullResult> PopMessageAsync(String brokerName, Int32 queueId = -1, Int32 maxNums = 32, Int64 invisibleTime = 60_000, Int32 pollTime = 15_000, CancellationToken cancellationToken = default)
14831484
{
14841485
if (String.IsNullOrEmpty(brokerName)) throw new ArgumentNullException(nameof(brokerName));
14851486

@@ -1493,6 +1494,7 @@ public async Task<PullResult> PopMessageAsync(String brokerName, Int32 maxNums =
14931494
{
14941495
consumerGroup = Group,
14951496
topic = Topic,
1497+
queueId,
14961498
maxMsgNums = maxNums,
14971499
invisibleTime,
14981500
pollTime,
@@ -1528,11 +1530,12 @@ public async Task<PullResult> PopMessageAsync(String brokerName, Int32 maxNums =
15281530

15291531
/// <summary>确认Pop消息消费完成</summary>
15301532
/// <param name="brokerName">Broker名称</param>
1531-
/// <param name="extraInfo">消息额外信息,Pop拉取时返回</param>
1532-
/// <param name="offset">消息偏移</param>
1533+
/// <param name="extraInfo">Pop检查点信息,即消息属性中的POP_CK字段值</param>
1534+
/// <param name="offset">消息在Queue中的偏移量</param>
1535+
/// <param name="queueId">队列编号</param>
15331536
/// <param name="cancellationToken">取消通知</param>
15341537
/// <returns></returns>
1535-
public async Task<Boolean> AckMessageAsync(String brokerName, String extraInfo, Int64 offset, CancellationToken cancellationToken = default)
1538+
public async Task<Boolean> AckMessageAsync(String brokerName, String extraInfo, Int64 offset, Int32 queueId = -1, CancellationToken cancellationToken = default)
15361539
{
15371540
using var span = Tracer?.NewSpan($"mq:{Name}:AckMessage", offset);
15381541
try
@@ -1546,6 +1549,7 @@ public async Task<Boolean> AckMessageAsync(String brokerName, String extraInfo,
15461549
topic = Topic,
15471550
extraInfo,
15481551
offset,
1552+
queueId,
15491553
};
15501554

15511555
await bk.InvokeAsync(RequestCode.ACK_MESSAGE, null, header, true, cancellationToken).ConfigureAwait(false);
@@ -1559,14 +1563,28 @@ public async Task<Boolean> AckMessageAsync(String brokerName, String extraInfo,
15591563
}
15601564
}
15611565

1566+
/// <summary>确认Pop消息消费完成。自动从消息属性中提取Pop检查点信息(POP_CK)</summary>
1567+
/// <param name="brokerName">Broker名称</param>
1568+
/// <param name="msg">通过Pop方式拉取的消息</param>
1569+
/// <param name="cancellationToken">取消通知</param>
1570+
/// <returns></returns>
1571+
public Task<Boolean> AckMessageAsync(String brokerName, MessageExt msg, CancellationToken cancellationToken = default)
1572+
{
1573+
if (msg == null) throw new ArgumentNullException(nameof(msg));
1574+
if (String.IsNullOrEmpty(msg.PopCheckPoint)) throw new ArgumentException("消息不含Pop检查点信息(POP_CK属性缺失),请确认该消息是通过Pop方式拉取的。", nameof(msg));
1575+
1576+
return AckMessageAsync(brokerName, msg.PopCheckPoint, msg.QueueOffset, msg.QueueId, cancellationToken);
1577+
}
1578+
15621579
/// <summary>修改Pop消息不可见时间,延长消费处理窗口</summary>
15631580
/// <param name="brokerName">Broker名称</param>
1564-
/// <param name="extraInfo">消息额外信息</param>
1565-
/// <param name="offset">消息偏移</param>
1581+
/// <param name="extraInfo">Pop检查点信息,即消息属性中的POP_CK字段值</param>
1582+
/// <param name="offset">消息在Queue中的偏移量</param>
15661583
/// <param name="invisibleTime">新的不可见时间(毫秒)</param>
1584+
/// <param name="queueId">队列编号</param>
15671585
/// <param name="cancellationToken">取消通知</param>
15681586
/// <returns></returns>
1569-
public async Task<Boolean> ChangeInvisibleTimeAsync(String brokerName, String extraInfo, Int64 offset, Int64 invisibleTime, CancellationToken cancellationToken = default)
1587+
public async Task<Boolean> ChangeInvisibleTimeAsync(String brokerName, String extraInfo, Int64 offset, Int64 invisibleTime, Int32 queueId = -1, CancellationToken cancellationToken = default)
15701588
{
15711589
using var span = Tracer?.NewSpan($"mq:{Name}:ChangeInvisibleTime", offset);
15721590
try
@@ -1581,6 +1599,7 @@ public async Task<Boolean> ChangeInvisibleTimeAsync(String brokerName, String ex
15811599
extraInfo,
15821600
offset,
15831601
invisibleTime,
1602+
queueId,
15841603
};
15851604

15861605
await bk.InvokeAsync(RequestCode.CHANGE_MESSAGE_INVISIBLETIME, null, header, true, cancellationToken).ConfigureAwait(false);
@@ -1594,6 +1613,20 @@ public async Task<Boolean> ChangeInvisibleTimeAsync(String brokerName, String ex
15941613
}
15951614
}
15961615

1616+
/// <summary>修改Pop消息不可见时间,延长消费处理窗口。自动从消息属性中提取Pop检查点信息(POP_CK)</summary>
1617+
/// <param name="brokerName">Broker名称</param>
1618+
/// <param name="msg">通过Pop方式拉取的消息</param>
1619+
/// <param name="invisibleTime">新的不可见时间(毫秒)</param>
1620+
/// <param name="cancellationToken">取消通知</param>
1621+
/// <returns></returns>
1622+
public Task<Boolean> ChangeInvisibleTimeAsync(String brokerName, MessageExt msg, Int64 invisibleTime, CancellationToken cancellationToken = default)
1623+
{
1624+
if (msg == null) throw new ArgumentNullException(nameof(msg));
1625+
if (String.IsNullOrEmpty(msg.PopCheckPoint)) throw new ArgumentException("消息不含Pop检查点信息(POP_CK属性缺失),请确认该消息是通过Pop方式拉取的。", nameof(msg));
1626+
1627+
return ChangeInvisibleTimeAsync(brokerName, msg.PopCheckPoint, msg.QueueOffset, invisibleTime, msg.QueueId, cancellationToken);
1628+
}
1629+
15971630
/// <summary>批量确认Pop消息消费完成</summary>
15981631
/// <param name="brokerName">Broker名称</param>
15991632
/// <param name="ackEntries">批量确认条目列表,每个条目包含extraInfo和offset</param>

NewLife.RocketMQ/Protocol/MessageExt.cs

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,13 @@ public class MessageExt : Message, IAccessor
5858

5959
/// <summary>消息编号</summary>
6060
public String MsgId { get; set; }
61+
62+
/// <summary>Pop检查点信息。Pop消费模式下由Broker在消息属性中返回,Ack/ChangeInvisibleTime操作时需传入此值</summary>
63+
public String PopCheckPoint
64+
{
65+
get => Properties.TryGetValue("POP_CK", out var str) ? str : null;
66+
set => Properties["POP_CK"] = value;
67+
}
6168
#endregion
6269

6370
#region 构造

XUnitTestRocketMQ/PopConsumeTests.cs

Lines changed: 107 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
using System;
22
using System.ComponentModel;
33
using NewLife.RocketMQ;
4+
using NewLife.RocketMQ.Protocol;
45
using Xunit;
56

67
namespace XUnitTestRocketMQ;
@@ -26,6 +27,16 @@ await Assert.ThrowsAsync<ArgumentNullException>(() =>
2627
consumer.PopMessageAsync(""));
2728
}
2829

30+
[Fact]
31+
[DisplayName("PopMessageAsync_可指定queueId参数")]
32+
public async void PopMessageAsync_WithQueueId_ThrowsWhenBrokerNameNull()
33+
{
34+
using var consumer = new Consumer();
35+
// 验证带queueId的重载依然会在brokerName为null时抛出异常
36+
await Assert.ThrowsAsync<ArgumentNullException>(() =>
37+
consumer.PopMessageAsync(null, queueId: 0));
38+
}
39+
2940
[Fact]
3041
[DisplayName("AckMessageAsync_无Broker连接时返回false")]
3142
public async void AckMessageAsync_NoBroker_ReturnsFalse()
@@ -36,6 +47,46 @@ public async void AckMessageAsync_NoBroker_ReturnsFalse()
3647
Assert.False(result);
3748
}
3849

50+
[Fact]
51+
[DisplayName("AckMessageAsync_指定queueId_无Broker连接时返回false")]
52+
public async void AckMessageAsync_WithQueueId_NoBroker_ReturnsFalse()
53+
{
54+
using var consumer = new Consumer();
55+
var result = await consumer.AckMessageAsync("nonexistent", "extra", 0, queueId: 2);
56+
Assert.False(result);
57+
}
58+
59+
[Fact]
60+
[DisplayName("AckMessageAsync_传入MessageExt_Null消息抛出异常")]
61+
public async void AckMessageAsync_NullMsg_ThrowsException()
62+
{
63+
using var consumer = new Consumer();
64+
await Assert.ThrowsAsync<ArgumentNullException>(() =>
65+
consumer.AckMessageAsync("broker", (MessageExt)null));
66+
}
67+
68+
[Fact]
69+
[DisplayName("AckMessageAsync_传入MessageExt_缺少POP_CK属性抛出异常")]
70+
public async void AckMessageAsync_MsgWithoutPopCk_ThrowsArgumentException()
71+
{
72+
using var consumer = new Consumer();
73+
var msg = new MessageExt { QueueId = 1, QueueOffset = 100 };
74+
// 没有设置 PopCheckPoint(POP_CK)
75+
await Assert.ThrowsAsync<ArgumentException>(() =>
76+
consumer.AckMessageAsync("broker", msg));
77+
}
78+
79+
[Fact]
80+
[DisplayName("AckMessageAsync_传入MessageExt_无Broker连接时返回false")]
81+
public async void AckMessageAsync_WithMsgExt_NoBroker_ReturnsFalse()
82+
{
83+
using var consumer = new Consumer();
84+
var msg = new MessageExt { QueueId = 1, QueueOffset = 100 };
85+
msg.PopCheckPoint = "100 1700000000000 60000 1 broker-a 1";
86+
var result = await consumer.AckMessageAsync("nonexistent", msg);
87+
Assert.False(result);
88+
}
89+
3990
[Fact]
4091
[DisplayName("ChangeInvisibleTimeAsync_无Broker连接时返回false")]
4192
public async void ChangeInvisibleTimeAsync_NoBroker_ReturnsFalse()
@@ -45,13 +96,65 @@ public async void ChangeInvisibleTimeAsync_NoBroker_ReturnsFalse()
4596
Assert.False(result);
4697
}
4798

99+
[Fact]
100+
[DisplayName("ChangeInvisibleTimeAsync_指定queueId_无Broker连接时返回false")]
101+
public async void ChangeInvisibleTimeAsync_WithQueueId_NoBroker_ReturnsFalse()
102+
{
103+
using var consumer = new Consumer();
104+
var result = await consumer.ChangeInvisibleTimeAsync("nonexistent", "extra", 0, 30000, queueId: 3);
105+
Assert.False(result);
106+
}
107+
108+
[Fact]
109+
[DisplayName("ChangeInvisibleTimeAsync_传入MessageExt_Null消息抛出异常")]
110+
public async void ChangeInvisibleTimeAsync_NullMsg_ThrowsException()
111+
{
112+
using var consumer = new Consumer();
113+
await Assert.ThrowsAsync<ArgumentNullException>(() =>
114+
consumer.ChangeInvisibleTimeAsync("broker", (MessageExt)null, 30000));
115+
}
116+
117+
[Fact]
118+
[DisplayName("ChangeInvisibleTimeAsync_传入MessageExt_缺少POP_CK属性抛出异常")]
119+
public async void ChangeInvisibleTimeAsync_MsgWithoutPopCk_ThrowsArgumentException()
120+
{
121+
using var consumer = new Consumer();
122+
var msg = new MessageExt { QueueId = 2, QueueOffset = 200 };
123+
// 没有设置 PopCheckPoint(POP_CK)
124+
await Assert.ThrowsAsync<ArgumentException>(() =>
125+
consumer.ChangeInvisibleTimeAsync("broker", msg, 30000));
126+
}
127+
128+
[Fact]
129+
[DisplayName("ChangeInvisibleTimeAsync_传入MessageExt_无Broker连接时返回false")]
130+
public async void ChangeInvisibleTimeAsync_WithMsgExt_NoBroker_ReturnsFalse()
131+
{
132+
using var consumer = new Consumer();
133+
var msg = new MessageExt { QueueId = 2, QueueOffset = 200 };
134+
msg.PopCheckPoint = "200 1700000000000 60000 1 broker-a 2";
135+
var result = await consumer.ChangeInvisibleTimeAsync("nonexistent", msg, 30000);
136+
Assert.False(result);
137+
}
138+
139+
[Fact]
140+
[DisplayName("MessageExt_PopCheckPoint属性读写正常")]
141+
public void MessageExt_PopCheckPoint_GetSet()
142+
{
143+
var msg = new MessageExt();
144+
Assert.Null(msg.PopCheckPoint);
145+
146+
msg.PopCheckPoint = "100 1700000000000 60000 1 broker-a 1";
147+
Assert.Equal("100 1700000000000 60000 1 broker-a 1", msg.PopCheckPoint);
148+
Assert.Equal("100 1700000000000 60000 1 broker-a 1", msg.Properties["POP_CK"]);
149+
}
150+
48151
[Fact]
49152
[DisplayName("RequestCode包含Pop消费相关码")]
50153
public void RequestCode_ContainsPopCodes()
51154
{
52-
Assert.Equal(200050, (Int32)NewLife.RocketMQ.Protocol.RequestCode.POP_MESSAGE);
53-
Assert.Equal(200051, (Int32)NewLife.RocketMQ.Protocol.RequestCode.ACK_MESSAGE);
54-
Assert.Equal(200052, (Int32)NewLife.RocketMQ.Protocol.RequestCode.CHANGE_MESSAGE_INVISIBLETIME);
55-
Assert.Equal(200151, (Int32)NewLife.RocketMQ.Protocol.RequestCode.BATCH_ACK_MESSAGE);
155+
Assert.Equal(200050, (Int32)RequestCode.POP_MESSAGE);
156+
Assert.Equal(200051, (Int32)RequestCode.ACK_MESSAGE);
157+
Assert.Equal(200052, (Int32)RequestCode.CHANGE_MESSAGE_INVISIBLETIME);
158+
Assert.Equal(200151, (Int32)RequestCode.BATCH_ACK_MESSAGE);
56159
}
57160
}

0 commit comments

Comments
 (0)