Skip to content

消息怎么不丢、不重?—— 生产消费语义 ​

属于 S3 Kafka 深入 · 第二篇 上一篇:架构与存储 下一篇:可靠性与积压

用消息队列,最头疼的两件事是"消息丢了"和"消息重复了"。这背后是生产者和消费者两端的语义设计。这一篇把它们拆开讲透,尤其把"幂等"和"rebalance"这两个抽象概念讲明白。

生产者:重试带来的重复,用幂等解决 ​

先看问题从哪来。生产者发消息失败会重试,但网络超时有个尴尬:消息可能已经写进 Broker 了,只是确认没传回来。重试就等于把同一条消息又写了一遍。这是"重复"的第一来源。

幂等生产者解决的就是"重试导致的重复"。它的机制很巧妙,用"快递单号"来理解:

每个生产者被分配一个 PID(Producer ID),相当于"寄件人编号";发给每个分区的消息还带一个单调递增的序列号(sequence),相当于"给这一家的包裹编号 1、2、3..."。Broker 记录每个分区已经收到的最大编号——收到编号 ≤ 已收到的,就认定是重复,直接丢弃。

这样就保证了"同一个生产者、同一个分区内"不会因重试而重复。但要记住幂等的边界:它只保证"单分区、单会话内"不重复。生产者重启后 PID 会变,就没法跨会话去重了——这正是"幂等生产者"和"事务"的分工所在。

事务:跨分区的原子性 ​

幂等只解决"单分区不重复",但如果一个业务要同时写多个分区,要求"要么全成功、要么全失败"呢?这就轮到事务了。

Kafka 的事务基于幂等实现,配合 transactional.id,通过一个"事务协调器"管理:开启事务 → 写入多个分区 → 一起提交(或中止)。事务 + 幂等 + 消费者的"只读已提交",才能拼出端到端的 exactly-once。

消费者:消费组和 offset 书签 ​

消费者这端,先理解两个概念。

消费组(Consumer Group) 是一组消费者协同消费一个 Topic。规则是:一个分区只能被组内一个消费者消费——像"一个座位只能坐一个人"。所以消费者实例数超过分区数时,多出的实例会空闲。这就是"分区数是并行度上限"的由来。

消费者的进度用 offset 记录,offset 就是"读到哪一页的书签"。这里有个关键选择:

  • 自动提交:定时自动挪书签,省事,但可能"书没看完书签就挪了"(处理前崩溃导致丢消息),或"书看完了书签没挪"(处理后崩溃导致重复)。
  • 手动提交:业务处理完再挪书签,配合"先处理后提交",实现 at-least-once(至少一次,可能重复,但绝不丢)。

三种投递语义要分清:at-most-once(最多一次,可能丢)、at-least-once(至少一次,可能重复,工程默认推荐)、exactly-once(恰好一次,需幂等 + 事务)。工程上普遍用 at-least-once + 消费者幂等,用幂等把重复消化掉。

rebalance:消费组的"重新分座位" ​

rebalance(再均衡) 是消费组里最容易踩坑的概念。想象一个班级(消费组)在换座位:有同学来了、有同学走了、或者座位(分区)数量变了,老师就要重新排座位——这就是 rebalance。

触发条件有三类:

  1. 消费者实例加入或退出(含宕机、心跳超时)。
  2. 订阅的 Topic 分区数变化。
  3. 消费者处理消息过久,超过 max.poll.interval.ms(默认 5 分钟)被踢出组。

最大的危害是:rebalance 期间,组内所有消费者停止消费——像换座位时全班都不许上课。频繁 rebalance 会导致吞吐骤降、消息大量积压。

怎么减少 rebalance:

  1. 调大 max.poll.interval.ms,避免处理慢被误踢。
  2. 调大 session.timeout.ms / heartbeat.interval.ms,避免网络抖动误判。
  3. 保持消费者实例数量稳定,别频繁上下线。
  4. 用静态成员(group.instance.id)——给每个消费者一个固定身份,重启后还坐原来的座位,不触发 rebalance。

串起来 ​

生产端的"不重"靠幂等(PID + sequence 防重试重复)和事务(跨分区原子);消费端的"不丢"靠手动提交 + at-least-once;而 rebalance 是消费组动态调整的代价,要理解它的触发条件(换座位)和规避手段(静态成员)。下一篇把这些拼成完整的可靠性方案,并处理"消息积压"这个线上常见故障。

下一篇讲可靠性与积压:把"不丢、不重、顺序"拼成完整方案,再解决消息积压。

持续学习,持续构建。