Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 12 additions & 1 deletion NewLife.RocketMQ/Consumer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -1170,14 +1170,25 @@ private async Task InitOffsetAsync(CancellationToken cancellationToken = default
{
offset = offsetTable.BrokerOffset;
if (offset <= 0) offset = await QueryMaxOffset(store.Queue, cancellationToken).ConfigureAwait(false);
// 主题尚未在Broker端创建时,QueryMaxOffset 返回 -1,从头开始消费
if (offset < 0) offset = 0;
}
else
{
offset = await QueryMinOffset(store.Queue, cancellationToken).ConfigureAwait(false);
if (offset < 0) offset = 0;
}

store.Offset = store.CommitOffset = offset;
await UpdateOffset(store.Queue, offset, cancellationToken).ConfigureAwait(false);
try
{
await UpdateOffset(store.Queue, offset, cancellationToken).ConfigureAwait(false);
}
catch (ResponseException ex) when (ex.Code == ResponseCode.TOPIC_NOT_EXIST)
{
// 主题尚未在Broker端创建,忽略偏移提交错误,首次消费后将自动提交偏移
WriteLog("主题[{0}]在Broker端尚未创建,跳过初始偏移提交", store.Queue.Topic ?? Topic ?? "unknown");
}
}
else
{
Expand Down
4 changes: 3 additions & 1 deletion NewLife.RocketMQ/MqBase.cs
Original file line number Diff line number Diff line change
Expand Up @@ -366,7 +366,9 @@ protected virtual void OnStart()

// 阻塞获取Broker地址,确保首次使用之前已经获取到Broker地址
var rs = client.GetRouteInfo(Topic);
DefaultTopicQueueNums = Math.Min(DefaultTopicQueueNums, rs.Where(e => e.Permission.HasFlag(Permissions.Write) && e.WriteQueueNums > 0).Select(e => e.WriteQueueNums).First());
var writeQueueNums = rs.Where(e => e.Permission.HasFlag(Permissions.Write) && e.WriteQueueNums > 0).Select(e => e.WriteQueueNums).FirstOrDefault();
if (writeQueueNums > 0)
DefaultTopicQueueNums = Math.Min(DefaultTopicQueueNums, writeQueueNums);

foreach (var item in rs)
{
Expand Down
28 changes: 15 additions & 13 deletions XUnitTestRocketMQ/Consumers/ConsumerRetryDLQIntegrationTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -105,23 +105,15 @@ public async Task ConsumerRetry_DisabledRetry_MessageNotRedelivered()
const String topic = "nx_no_retry_test";
const String group = "nx_no_retry_group";

using var producer = new Producer
{
Topic = topic,
NameServerAddress = set.NameServer,
Log = XTrace.Log,
};
producer.Start();
Thread.Sleep(2000);

var stamp = DateTime.UtcNow.Ticks.ToString();
producer.Publish($"no-retry-{stamp}");
Thread.Sleep(500);

var receiveCount = 0;
var receivedOnce = new SemaphoreSlim(0, 1);
var advancedOffset = new SemaphoreSlim(0, 1);

// 提前生成唯一标识,避免 lambda 闭包引用顺序问题
var stamp = Guid.NewGuid().ToString("N")[..8];
Comment on lines +112 to +113

// 先启动消费者并等待重新平衡,确保消费者就绪后再发送消息
// 若消息在消费者启动前发送,FromLastOffset=true 会跳过该消息
using var consumer = new Consumer
{
Topic = topic,
Expand Down Expand Up @@ -158,6 +150,16 @@ public async Task ConsumerRetry_DisabledRetry_MessageNotRedelivered()
return true; // 第二次起放行,偏移推进
};
consumer.Start();
Thread.Sleep(3000); // 等待消费者重新平衡完成

using var producer = new Producer
{
Topic = topic,
NameServerAddress = set.NameServer,
Log = XTrace.Log,
};
producer.Start();
producer.Publish($"no-retry-{stamp}");

// 等待收到一次
var got = await receivedOnce.WaitAsync(TimeSpan.FromSeconds(30));
Expand Down