NATS JetStream实现Exactly Once Delivery消费者端方案咨询
NATS JetStream 端到端精确一次投递(EOS)消费者侧落地指南
你已经掌握的生产侧MsgId+流级别重复检测窗口的配置,已经覆盖了生产侧的幂等发消息能力,要补全端到端EOS,核心要解决的是消费者侧Ack丢失、服务宕机导致的重复投递、消息丢失问题——JetStream本身的重投逻辑是消息在AckWait窗口内没收到确认就会重投,本身只保证At Least Once语义,消费者侧的配置和逻辑要把投递语义抬升到Exactly Once,具体做法如下:
一、消费者基础配置项(必须配置)
- 强制关闭所有SDK提供的自动Ack、异步批量Ack选项,必须用手动显式Ack模式,绝对不能拿到消息就先回Ack再处理业务。
- 创建消费者时必须指定固定的
Durable持久化名称,不要用临时Ephemeral消费者,避免消费者重启后消费位点错乱拉取重复历史消息。 - Ack策略固定选
ExplicitAck(大部分SDK默认是这个,手动确认单条消息),不要选AckAll、None这类批量确认或者无确认策略。 AckWait(Ack等待超时)设置为单条消息最大业务处理耗时的1.3~1.5倍,留足网络和GC余量,避免业务还没处理完Broker就触发重投。MaxDeliver(最大投递次数)不要设为无限,建议设3~5次,超过次数的消息投递到配套死信流,避免异常消息卡死整个消费链路。- 并发度、单次Pull拉取量要和业务处理能力匹配,不要拉取超过
AckWait时间内能处理完的消息量,避免消息堆在客户端内存里服务宕机导致丢失。
二、消费端核心业务逻辑(无代码逻辑仅配参数无法实现EOS)
所有逻辑围绕*「业务处理成功」和「消息确认」的原子性*设计,没有任何Broker侧配置能替你完成业务层的幂等校验:
- 拿到消息第一时间先提取消息头里的
Nats-Msg-Id(就是生产者侧设置的去重ID),先查幂等记录,再执行业务。注意:幂等记录必须和你的业务库存在同一个事务上下文里,不要单独用Redis、其他独立存储做幂等,避免分布式事务不一致。比如你用MySQL存业务数据,就单独建一张消费幂等表,字段存
msg_id、process_status、create_time,和业务表在同一个库实例里。 - 如果幂等表查到这个
msg_id已经是处理成功状态,直接给Broker回Ack跳过这条消息,不要执行业务逻辑。 - 如果幂等表没查到记录,开启本地数据库事务:先执行业务逻辑的增删改操作,再往幂等表里插入当前
msg_id对应的处理成功记录,两个操作在同一个事务里提交。 - 只有等数据库事务100%提交成功之后,再给JetStream回对应消息的Ack。如果事务提交失败、业务处理报错,直接回Nack或者不回Ack,等Broker超时重投就行,绝对不要在事务提交前发Ack。
逻辑兜底原理
- 如果事务提交成功,哪怕回Ack的过程中网络断了、服务宕机了,Broker超时重投同一条消息,下次消费时第一步查幂等表就会发现已经处理过,直接Ack跳过,不会重复执行业务。
- 如果事务没提交成功服务就宕机了,消息不会收到Ack,Broker会按规则重投,不会出现消息丢了没处理的情况。
三、常见误区避坑
- 不要觉得开了JetStream的事务消息就能实现消费侧EOS:JetStream的事务能力只覆盖生产侧跨多流原子发消息的场景,管不到消费侧的业务重复执行问题。
- 不要用批量Ack跨消息确认:比如一次拉10条消息,处理到第5条服务崩了,要是前4条已经提前发了Ack但业务事务没提交,就会直接丢消息。
- 不要把幂等校验放在业务逻辑执行之后:如果业务逻辑有外部调用(比如发通知、扣减库存),没前置幂等的话重投就会导致重复调用外部接口。
内容的提问来源于stack exchange,提问作者bulwark
相关产品推荐
相关产品推荐

