You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

.NET 8中NATS JetStream消费者重复处理已确认消息问题

问题分析与解决方案

你遇到的核心问题是:设置MaxDeliver后,即便调用了msg.AckAsync()确认消息,每条消息仍会被重复处理MaxDeliver次;不设置该参数时则正常。这是因为持久化消费者的递送策略与MaxDeliver的交互逻辑,加上代码中每次运行都会重置消费者配置导致的。

具体原因

  1. DeliverPolicy=All的影响:你设置了DeliverPolicy.All,意味着消费者会从流的起始位置消费所有消息,包括历史消息。如果每次运行程序都重新创建/更新消费者,它会重复拉取所有已存在的消息。
  2. MaxDeliver的误解:该参数定义的是消息未被确认时的最大重试递送次数(包含首次递送),但如果消费者每次启动都重新拉取所有消息,会被JetStream判定为新的递送尝试,从而触发重复处理。
  3. 持久化消费者的状态残留:你给消费者设置了固定名称first-new-consumer,属于持久化消费者,其消费状态会被NATS服务器保存。每次运行程序调用CreateOrUpdateConsumerAsync时,会重置或沿用旧状态,导致重复消费。

修复方案

1. 调整消费者递送策略

将DeliverPolicy改为New,让消费者仅处理创建之后发布的新消息,避免重复消费历史数据:

ConsumerConfig consumerConfig = new ConsumerConfig(name: "first-new-consumer")
{
    DeliverPolicy = ConsumerConfigDeliverPolicy.New, // 只消费新消息
    MaxDeliver = 2,
    AckWait = TimeSpan.FromSeconds(30), // 显式设置确认超时时间,确保Ack能被服务器接收
    // Backoff = backOff // 可先注释,排除退避策略的干扰
};

2. 避免重复创建/更新消费者

添加消费者存在性检查,仅在消费者不存在时创建,避免每次运行重置配置:

var consumerExists = false;
try
{
    await js.GetConsumerAsync(streamName, "first-new-consumer");
    consumerExists = true;
}
catch (NatsJSException)
{
    // 消费者不存在,捕获异常后创建
}

var consumer = consumerExists 
    ? await js.GetConsumerAsync(streamName, "first-new-consumer") 
    : await js.CreateConsumerAsync(streamName, consumerConfig);

3. 验证Ack是否生效

在AckAsync()后添加日志,确认确认操作完成,同时使用NATS CLI工具检查流和消费者状态:

# 查看流信息,确认消息确认情况
nats stream info sample-new

# 查看消费者状态,检查递送计数
nats consumer info sample-new first-new-consumer

4. 使用临时消费者测试(可选)

如果不需要持久化消费状态,可以创建临时消费者(不设置name),避免状态残留问题:

var consumerConfig = new ConsumerConfig
{
    DeliverPolicy = ConsumerConfigDeliverPolicy.New,
    MaxDeliver = 2
};
var consumer = await js.CreateConsumerAsync(streamName, consumerConfig);

修改后的完整代码示例

using NATS.Client.Core;
using NATS.Client.JetStream;
using NATS.Client.JetStream.Models;
using NATS.Client.Serializers.Json;
using System.Text.Json;

try
{
    int count = 0;
    var natsOpts = NatsOpts.Default with
    {
        SerializerRegistry = NatsJsonSerializerRegistry.Default,
        Url = "nats://localhost:53810",
        AuthOpts = new NatsAuthOpts()
    };

    string subject = "test.new.*";
    string streamName = "sample-new";
    string consumerName = "first-new-consumer";

    await using var nats = new NatsConnection(natsOpts);
    var js = new NatsJSContext(nats);

    // 创建流(仅当不存在时)
    try
    {
        await js.CreateStreamAsync(new StreamConfig(name: streamName, subjects: new[] { subject }));
    }
    catch (NatsJSException)
    {
        // 流已存在,忽略
    }

    // 发布新消息
    for (var i = 0; i < 10; i++)
    { 
        var message = new SampleEvent("John.Doe", "Sample Message", DateTime.Now.ToString(), i.ToString());
        var ack = await js.PublishAsync("test.new.first", message);
        ack.EnsureSuccess();
    }

    // 配置消费者
    var consumerConfig = new ConsumerConfig(name: consumerName)
    {
        DeliverPolicy = ConsumerConfigDeliverPolicy.New,
        MaxDeliver = 2,
        AckWait = TimeSpan.FromSeconds(30)
    };

    // 检查消费者是否存在,不存在则创建
    var consumer = await js.CreateOrUpdateConsumerAsync(streamName, consumerConfig);

    // 消费消息
    await foreach (var msg in consumer.ConsumeAsync<SampleEvent>())
    {
        var order = msg.Data;
        Console.WriteLine($"Processing {order.ToString()}...");
        await msg.AckAsync();
        Console.WriteLine($"Processed count: {++count}"); 
    }
}
catch (Exception ex)
{
    Console.WriteLine(ex.Message);
}

public class SampleEvent
{
    public SampleEvent() {}

    public SampleEvent(string userName, string message, string utcTime, string count)
    {
        UserName = userName;
        Message = message;
        UtcTime = utcTime;
        Count = count;
    }

    public string UserName { get; set; }
    public string Message { get; set; }
    public string UtcTime { get; set; }
    public string Count { get; set; }

    public override string ToString() => JsonSerializer.Serialize(this);
}

内容的提问来源于stack exchange,提问作者Kunal

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 09:44:53