Skip to content

消息不丢、不重、不乱序,还积压了怎么办?—— 可靠性与积压 ​

属于 S3 Kafka 深入 · 第三篇 上一篇:生产消费语义 下一篇:面试题集

上一篇讲了生产者、消费者各自的语义。这一篇把"不丢、不重、顺序"拼成完整方案,再处理线上最头疼的"消息积压"。这三个"不",每一个背后都是一整套机制。

不丢:三道闸门都要设防 ​

消息从生产到消费,要过"生产者 → Broker → 消费者"三道闸门,任何一道疏忽都会漏。

第一道:生产端。 要 acks=-1(所有 ISR 副本都确认才返回)+ retries>0(失败重试)+ min.insync.replicas>=2。后一个参数容易被忽略——上一篇说过,acks=-1 在 ISR 只剩 Leader 时会退化成 1,而 min.insync.replicas>=2 就是防止这种退化:如果同步副本不够 2 个,直接拒绝写入,而不是悄悄降级。

第二道:Broker 端。 副本数 >=3,unclean.leader.election.enable=false。后者保证 Leader 挂了只从 ISR 里选新 Leader,宁可不选也不选一个落后的副本(选落后副本 = 必然丢数据)。

第三道:消费端。 手动提交 offset,且坚持"先处理业务、再提交"。这个顺序是灵魂:如果先提交再处理,处理时崩溃,这条消息就"被书签跳过"了,永久丢失。

不重:靠消费端幂等消化 ​

at-least-once 语义下,重复是常态。别指望"让消息恰好只投递一次",而是要让重复消费的结果幂等——重复消费一百次,结果和消费一次一样。三种常见手段:

  1. 唯一键/唯一索引:消息带业务唯一 ID(如订单号),落库时靠数据库唯一索引去重,重复插入直接报错。这是最可靠的兜底。
  2. 去重表:单独一张表记录"已处理的消息 ID",消费前查一下,处理过就跳过。
  3. 状态机:更新前校验状态,比如订单只能从"待支付"变"已支付"一次,重复的更新请求天然失效。

用比喻记:唯一键是"身份证号查重",去重表是"已办事项清单",状态机是"流程只允许单向走"。

顺序:只保证单分区有序 ​

Kafka 只保证单个分区内有序,不保证跨分区有序。这像流水线:同一条流水线(分区)上,东西是按顺序加工的;但不同流水线之间,谁先谁后说不准。

要保证顺序,本质是"让需要有序的消息进同一个分区":

  • 全局有序 → 单分区(牺牲并行,吞吐低,极少用)。
  • 同 key 有序 → 按 key 哈希到同一分区(最常用,如同一订单的所有消息都进同一分区)。

消费端还有个容易被忽略的点:一个分区只被一个消费者消费,但消费者内部如果多线程并发处理,还是可能乱序。所以要分区级串行——每个分区对应一个处理协程,保证单分区内的处理顺序。

积压:消息堵车了 ​

消息积压的典型信号是消费组的 lag(落后量,即"生产到了第 100 万条,消费才到第 50 万条")持续增大。用交通比喻,lag 就是"堵车的车队长龙"。

堵车的原因不外乎:消费者实例太少(车道不够)、消费逻辑太慢(每辆车过收费站都磨蹭)、频繁 rebalance(动不动就重新排座位)、流量洪峰(突然涌来大量车)。

处理步骤(按顺序):

  1. 看 lag 定位:kafka-consumer-groups --describe 看落后量,先搞清楚堵在哪。
  2. 临时扩容消费者:加消费者实例——但记住扩容上限 = 分区数,实例超过分区数也没用(一个分区只能一个人消费,多出的实例只能干看着)。
  3. 优化消费逻辑:批量拉取、异步处理、去掉重复查库。
  4. 极端兜底:临时建一个新 topic,把积压的消息转储过去,用大量消费者并行消化,处理完再回写——相当于"把堵的车先引流到临时停车场分批疏散"。

长期看,要增加分区、优化下游、做好容量评估。


串起来 ​

可靠性的三件事——不丢(三道闸门设防)、不重(幂等消化)、顺序(同 key 同分区 + 分区串行)——拼起来就是一条可靠的链路;积压是这条链路超载的信号,按"定位 → 扩容 → 优化 → 兜底"的顺序处理,并且始终记住那个关键约束:消费者扩容上限 = 分区数。

下一篇是 S3 的面试题集,把 Kafka 的高频问题收口,练到连续追问三层。

持续学习,持续构建。