做分布式系统,消息投递是绕不开的坎。你肯定听说过三种语义:至多一次、至少一次、精确一次。但大部分时候,我们用的其实是“至少一次”。不是因为不想精确,是因为实现精确一次的代价太高,高到有时候得不偿失。可“至少一次”就那么简单吗?天真了,兄弟。这玩意儿表面上是个重试机制,底层却藏着大量的协议博弈和状态机陷阱。今天我们不聊虚的,就直接把“至少一次”扒个精光。从TCP的滑动窗口到Kafka的幂等事务,从理论上的美好到压测下的血泪史,一条条掰扯清楚。
说起来你可能不信,我第一次在生产环境被“至少一次”教育,是在一个促销活动的晚上。订单推送,消费者那边老是抱怨收到重复推送,老板盯着大屏,我汗流浃背。查了半天,就因为是消息中间件的默认投递语义就是至少一次,而消费者没做幂等。哎,多么痛的领悟。所以,别等到故障了才去补课。
那“至少一次”到底是怎么做到的?说白了,就是确认机制加无限重试。发送方发出一条消息,必须等到接收方的确认(ACK)才会认为成功;如果超时没收到ACK,就重发。接收方呢,处理完消息,发ACK。看上去完美,对吧?
错。这是理想模型。现实是——网络会抖动,节点会宕机,TCP连接会断开。任何一个环节出问题,都会导致重试。而重试就意味着消息可能重复。但为什么还是这么多系统用这个语义?因为它在可靠性和复杂度之间找到了一个诡异的平衡点。实现精确一次需要分布式事务或者额外的去重机制,成本太高,而很多业务能容忍少量重复,只要做好幂等。
但话说回来,就算你做好了幂等,“至少一次”带来的问题可远不止重复。乱序、死信、性能抖动,这些雷你踩过几个?
从源头说起:消息的“投”与“达”
先搞清楚一个概念:消息投递的承诺,是在哪个层面做出的?很多人把TCP的可靠传输和消息中间件的投递混为一谈。TCP保证的是数据包的可靠传输,但那是在连接之内。一旦连接断开,TCP就不管了。而消息中间件的“至少一次”,是业务层面的可靠性,它建立在不可靠的网络上。这就要求中间件在应用层实现一套确认和重传机制。
我们以RabbitMQ为例——虽然Kafka现在很火,但它的底层逻辑其实类似。生产者发送消息到Exchange,消息会被路由到队列。消费者从队列拉取或由Broker推送消息。关键点来了:消费者处理完消息后,必须显式发送ACK。如果消费者在发送ACK之前挂了或者连接断了,Broker就会认为消息未送达,重新把消息放回队列并投递给其他消费者(或者等这个消费者重连后再次投递)。这就是“至少一次”的根苗。
但这是不是绝对的?当然不是。如果Broker在收到ACK后宕机了,但ACK还没来得及落盘,恢复后消息状态可能还是未确认,又会导致重投。所以,持久化和确认机制的配合就至关重要。RabbitMQ提供了事务和confirm模式,Kafka则依赖ISR和Leader选举。
说到Kafka,它实现“至少一次”的方式更变态一些。消费者通过offset标记消费进度。如果消费者处理完消息,提交了offset,然后宕机,新启动的消费者会从上次提交的offset继续消费,不会重复。但如果先处理消息,提交offset之前宕机,或者offset提交失败,消息就会被重复消费。Kafka默认也是“至少一次”。要想精确一次,得开启幂等生产者和事务,代价嘛——吞吐量下降20%到30%,甚至更多。鱼和熊掌。

所以,不要神话“至少一次”。它就是一个妥协方案,在保证消息不丢的前提下,允许重复。那能不能做到不重复?能,但需要额外机制,后面再骂。
协议与状态机的博弈:ACK风暴与重试雪崩
“至少一次”不是孤立存在的。它和重试策略、超时设置、ACK机制深度耦合。这里面的学问,比你想的深。比如一个简单的问题:重试几次?间隔多久?
很多开发者喜欢用固定间隔重试3次。天真的做法。一旦依赖的服务出现短暂故障,3次固定间隔重试很可能全部失败,然后消息就被丢进死信队列了。但你如果用指数退避重试,消费者的处理延迟就会变得很不稳定。更要命的是,如果消费者的下游也依赖“至少一次”,就会形成重试链式反应,最后整个系统被重试流量打崩——俗称“重试风暴”。
我见过最惨的案例:一个物流系统在双十一流量下,因为消费者处理慢,消息积压,触发大量重试;重试又加重下游数据库压力,导致处理更慢,接着更多重试……死循环。最后数据库挂了,全链路阻塞。复盘一看,就是因为重试策略和超时时间没有根据业务量级和依赖的SLA闭环设计。
再有就是ACK的代价。在RabbitMQ里,每一条消息的ACK都是一个网络往返。如果用自动ACK,吞吐量是高,但那是“至多一次”,消息根本没保障。如果每条都手动确认,吞吐量直接腰斩。所以通常的做法是批量确认:处理好几条,一次性ACK。但批量确认有风险——如果消费者崩溃,这批消息都得重发,重复量更大。又是一个权衡。
而Kafka里的offset提交也是一样。自动提交间隔太长,重复量大;间隔太短,性能差。调优到一个合理值,需要根据消息处理耗时和可接受的重复窗口来算。比如你允许最多2秒的重复,那么在提交间隔加处理时间加网络往返不能超过2秒。这些细致活,没做过压测根本拿不准。

更隐蔽的坑:乱序。因为重试,消息的顺序很可能被打乱。消费者收到消息的次序和生产者发送的次序不一致。在分布式分区下,即使单一分区有序,重试也会打破顺序。如果要保证顺序,就得牺牲并行度——单线程消费,或者在业务层做顺序重排。这也是“至少一次”的代价。
数据的代价:压测下的赤裸现实

空口无凭,上数据。我们做过一轮针对Kafka“至少一次”与“精确一次”的对比压测。环境:3个Broker,1个Topic 3分区,异步生产,消费端单分区单线程,处理逻辑只做简单JSON解析并写入MySQL。消息大小1KB。压测时长30分钟。
结果:
“至少一次”模式(默认,手动提交offset,提交间隔100ms):平均吞吐量 125,000 msg/s,端到端延迟 P99 为 350ms,消息重复率 0.003%(因为少量消费者重启导致offset未及时提交)。
“精确一次”模式(开启幂等和事务,生产者事务API,消费者隔离级别read_committed):平均吞吐量 89,000 msg/s,P99 延迟 780ms,重复率 0%。
看到没,吞吐量下降近30%,延迟翻倍。这还是只做简单处理。如果业务逻辑复杂,事务开销占比较小,差异可能缩小,但无论如何,精确一次的性能损失是实实在在的。所以,如果你的业务能容忍少量重复,并且做好了幂等,“至少一次”是更经济的选择。
但重复率真的可控吗?上面的0.003%是在消费者平稳运行下的数据。如果我们刻意模拟故障:每5分钟kill掉一个消费者进程,重复率飙升到0.5%。500条里就有2-3条重复。如果消息量是每天几十亿,那就是个不能忽视的数字。所以,你要么接受它并做好业务去重,要么花大价钱上精确一次。
另一个压测是关于重试风暴。我们用RabbitMQ,100个消费者处理消息,下游依赖一个HTTP服务。故意让HTTP服务响应时间从50ms逐渐增加到5s。在无流控的情况下,消息积压导致连接数暴涨,重试消息大量涌入,最终RabbitMQ内存打满,节点重启。之后加上消费者端限流(prefetch count设为50),并且设置消息TTL和死信队列,再跑同样的场景,系统只是延迟变大,没有崩溃。这就是限流和死信的作用。
落地三大坑,以及怎么爬出来

好了,理论讲够了,直接上干货。根据多年踩坑经验,我把“至少一次”落地最容易掉进去的坑列出来,附上解决方案。
坑1:幂等做得太晚
很多团队知道要幂等,但都是在业务逻辑层做的,比如数据库唯一键。这没错,但不够。因为重试可能会导致你的服务对外部资源重复操作,比如发短信、扣款。如果扣款请求先到了银行,你后来才在数据库发现重复并回滚,银行那边可不会认。所以,幂等必须在入口处实现。最简单的办法:为每条消息生成一个全局唯一的业务ID(比如订单号 + 消息类型 + 时间戳),消费者收到消息后,先查这个ID是否处理过。可以用Redis的SETNX,或者数据库的去重表。处理成功后,再记下ID。注意,查重和存ID的操作必须是原子的,否则有并发问题。Redis的SETNX是原子的,数据库可以用唯一索引加INSERT ON DUPLICATE KEY。
坑2:死信队列变成垃圾场
重试耗尽后,消息进入死信队列。然后呢?很多团队就让它烂在里面。等到发现数据不一致,死信队列已经堆了几十万条消息,处理起来痛苦不堪。必须建立死信队列的监控和自动回收机制。我的做法:对每条死信做分类,能自动重试的(比如临时网络错误)用脚本定期重放;不能自动重试的(比如数据格式错误)落库并告警,人工介入。同时,在死信处理完之前,业务不能停止,就要设计补偿表,定期比对差异。
坑3:乱序导致的状态错乱
对于状态机流转的业务,顺序就是命。比如订单从“已支付”到“已发货”,如果“已发货”消息先到,“已支付”后到,消费者如果无脑处理,就会出一个“已支付”的订单,但发货行为已经发生。解决之道:要么严格让同一单据的消息走同一分区(用业务key哈希),然后在消费端单线程顺序处理;要么在消息体里带上版本号或时序字段,消费者维护一个状态缓存,忽略掉旧版本的消息。单线程消费会限制吞吐量,但可以通过分区扩容来线性扩展,这就是Kafka的设计精髓。
最后一个坑点我额外加一个:超时风暴与雪崩保护。消费者处理慢,导致重试,重试加重负载,恶性循环。除了前面说的消费者端限流,还必须设置断路器。当消息处理失败率超过阈值时,熔断这条消费链路,直接快速失败,让重试消息进入死信,等依赖恢复再慢启动。熔断的方案可以用Hystrix或Resilience4j,不赘述。
写完这些,你是不是觉得“至少一次”没那么简单?其实每一种技术选型背后都是权衡。你嫌弃它的重复,却又离不开它的可靠。工程就是这样,没有银弹,只有最适合的折衷。下次有人忽悠你说,我们的消息系统绝对不丢不重,你可以微微一笑:那你的吞吐量和延迟是多少?