Kafka(下):生产端、broker 与消费端的可靠性配置

📅
2 分钟阅读
·

系列目录

  1. MySQL 索引与慢查询:B+ 树如何减少扫描
  2. 声明式事务之下:InnoDB 的 MVCC 与锁
  3. 接入层:Nginx 反向代理与 OpenResty 的边界
  4. Web 容器与 Netty:线程模型之下的 IO 模型
  5. Redis(上):缓存用法与单线程模型的限制
  6. Redis(下):超出缓存用途的用法:锁、队列与排行榜
  7. 数据库访问层:连接池与 MyBatis 的显式 SQL
  8. Kafka(上):吞吐的来源是顺序 IO
  9. Kafka(下):生产端、broker 与消费端的可靠性配置(本篇)

一次位移提交导致的漏处理

本篇讨论 Kafka 在什么条件下不会丢失消息。(Kafka 异步化与削峰)提到当时接入 Kafka 时,消费端用去重表实现幂等,并通过定时对账处理遗漏,业务侧已经具备处理重复消息的机制。当时没有逐项核对生产端、broker 和消费端的配置,直接把「Kafka 保证不丢消息」作为前提。2018 年的一次对账发现了这个问题。

一个消费组的处理耗时偏长,enable.auto.commit=true(默认),auto.commit.interval.ms 使用默认值。某次对账发现,这个组在一段时间窗口内的消息计数低于生产侧,对应的消息未被处理。

复盘定位到的机制:自动提交提交的是「当前位点」,即最近一次 poll() 返回的最后一条消息之后的位置。poll() 把消息返回给应用时,位点就已经推进到这批末尾,先于应用处理。处理慢的情况下,周期性提交把整批位移提交上去时,这批消息可能还没处理完。这批里如果有的消息处理失败被异常捕获吞掉,或者消费者在处理中途重启,下次消费从已提交的位点继续,失败的那几条就落在已提交位点之后,不会被重新拉取,等于被跳过。投递语义从「至少一次」退化成了「至多一次」。

跳过的消息数量与单批消息数量相当(数字脱敏)。临时处理按时间戳回溯缺口区间重新消费,并依靠已有的去重表避免重复处理。复盘后关闭 enable.auto.commit,改为处理完一批再手动提交。该事件说明,投递语义取决于生产端、broker 和消费端的一组配置。2016 系列第 6 篇说明了消费幂等的约定(去重表、对账);本篇说明 Kafka 配置如何影响投递语义,不重复展开业务侧处理。

生产端、broker 与消费端的可靠性配置

Kafka 的可靠性配置分布在生产端、broker 副本和消费端。任一环节的配置不满足要求,消息可能丢失、重复或被跳过。

Kafka 可靠性配置:生产端、副本与消费端

生产端:acks 决定消息持久化到几个副本

生产端 acks 参数有三个取值:

  • acks=0:producer 发送后即视为成功,不等待 broker 确认。网络丢包或 broker 未收到消息时,producer 无法感知消息丢失。
  • acks=1:leader 写入后确认。leader 在确认后、follower 尚未同步时宕机,消息会丢失。
  • acks=all(或 -1):leader 等 ISR 里所有副本都写入后才确认。确认前消息已写入当前 ISR 中的所有副本。

ISR(in-sync replicas)是当前能跟上 leader 的副本集合。它不等于配置的副本总数。replica.lag.time.max.ms 决定 follower 多久未追上 leader 的最新位移会被移出 ISR。ISR 会动态变化,副本数配置为 3 时,ISR 中的副本数仍可能少于 3。

acks=all 的效果取决于 ISR 中的副本数。ISR 只剩 leader 时,acks=all 的确认条件与 acks=1 相同。min.insync.replicas 要求 ISR 中至少有指定数量的副本;数量不足时,producer 写入失败并抛出 NotEnoughReplicasExceptionmin.insync.replicas≥2 配合 acks=all 时,producer 只会在消息至少写入两个副本后收到确认。ISR 不足会拒绝写入。常见配置为副本数 3、min.insync.replicas=2acks=all:一个副本宕机后,ISR 仍有 2 个副本,写入继续;两个副本同时宕机后,ISR 只剩 1 个副本,写入失败。

broker 副本:HW、LEO 与 unclean leader election

LEO(log end offset)是每个副本日志的下一条写入位置;HW(high watermark)是所有 ISR 副本中最小的 LEO,即消息已被全部 ISR 副本确认的位置。消费者只能读到 HW 以下的消息。这个设计保证了消费者读到的消息一定已同步到全部 ISR 副本,leader 切换后不会出现「消费者读过但 follower 没有」的情况。

leader 宕机时,controller 从 ISR 中选出新 leader。ISR 中存在副本时,新 leader 包含完整的 HW 以下数据。unclean leader election 决定 ISR 为空时(全部副本滞后或宕机),是否允许从非 ISR 副本中选出 leader。

  • 打开(unclean.leader.election.enable=true):集群保持可用,但新 leader 的数据落后于旧 leader,旧 leader 上已确认、还没同步到新 leader 的消息丢失。
  • 关闭(false):ISR 为空时拒绝选举,该分区不可用,直到原 leader 恢复或 ISR 里有副本重新跟上。

该配置在分区可用性和已确认消息保留之间作出选择。日志类业务能够接受少量消息丢失且停机代价更高时,可以开启;账务类业务不能丢失已确认的消息时,应关闭。该参数的默认值在 0.11 从 true 改为 false。通知和统计链路配置为 false,以接受短暂不可用来保留已确认消息;确实需要优先保持可用的日志管道单独配置为 true。

消费端:位移提交时机与投递语义

消费端位移(offset)记录消费组在每个分区上消费到哪了。提交时机决定投递语义。

enable.auto.commit=true(默认)时,消费者按 auto.commit.interval.ms 周期性提交当前位点。当前位点是 poll() 返回给应用的最后一条消息之后的位置,poll() 返回消息时位点就已推进,先于应用处理。提交的位移只表示消息已拉取,不表示消息已处理。处理慢或处理失败时,位移已经推进到未处理消息之后,下次消费从已提交位点继续,未处理的部分被跳过,投递退化成至多一次。开头的漏处理由这一机制导致。

关闭自动提交后,手动提交可采用三种顺序:

  1. 先提交后处理poll() 返回消息后先 commitSync(),再处理。处理失败时位移已提交,消息丢失,投递语义为至多一次。
  2. 先处理后提交:处理完这批消息再 commitSync()。处理完、提交前进程崩溃时,重启后会重新消费这批消息,投递语义为至少一次,并会产生重复消息。该顺序适用于消费端能够幂等处理重复消息的场景。
  3. 处理与提交原子化:把消息处理和位移提交放进同一个事务。Kafka 0.11 引入了事务和精确一次语义,适用范围是「消费 Kafka、处理、写回 Kafka」的链路。处理结果写入外部系统(数据库、下游接口)时,事务无法将外部写入和位移提交原子化。0.11 当时刚发布,团队没有采用。多数场景采用第 2 种顺序,并通过幂等处理重复消息。

commitSync 同步阻塞,提交成功后才返回,会影响吞吐。commitAsync 异步不阻塞,但提交可能失败且不会自动重试。实践中可在循环中使用 commitAsync 减少阻塞,在结束时用 commitSync 提交最后一批位移。

rebalance:慢消费可能触发消费停顿和重复处理

consumer group 的 rebalance 指分区分配关系重新调整。触发条件:消费者加入或离开 group(部署、扩容、消费者崩溃)、订阅的分区数变化。rebalance 期间所有消费者停止消费,等待重新分配完成,这段时间消费停顿。

慢消费会增加 rebalance 的发生概率。max.poll.interval.ms 控制两次 poll() 之间的最大间隔:消费者处理一批消息的时间超过这个值后,broker 将该消费者视为失效并移出 group,触发 rebalance。rebalance 后分区分配给其他消费者,原消费者处理完该批消息后可能因已不在 group 中而提交失败;新消费者从上次提交的位移开始消费,可能重复处理。

处理时间过长会超过 max.poll.interval.ms,触发 rebalance 和消费停顿;停顿会增加积压,积压可能进一步延长处理时间。max.poll.interval.ms 应大于最坏处理时间,max.poll.records 用于限制单批拉取量。调参只能让处理时间保持在 rebalance 阈值内,慢消费的根因仍需由业务处理逻辑解决。

事件后的客户端配置

事件后,团队采用以下客户端配置,以实现至少一次投递并由业务侧处理重复消息(脱敏):

// producer:确认前持久化到全部 ISR 副本
Properties p = new Properties();
p.put("bootstrap.servers", "broker1:9092,broker2:9092,broker3:9092");
p.put("acks", "all");
p.put("retries", Integer.MAX_VALUE);   // 网络抖动时重试
// 无幂等生产者时,保证重试不打乱分区内顺序
p.put("max.in.flight.requests.per.connection", "1");
// 0.11 幂等生产者当时刚出,评估后未上;若开启可放宽上面的限制到 5

// consumer:处理后提交,避免位移越过未处理消息
Properties c = new Properties();
c.put("bootstrap.servers", "broker1:9092,broker2:9092,broker3:9092");
c.put("enable.auto.commit", "false");          // 关掉自动提交
c.put("max.poll.interval.ms", "300000");       // 大于最坏处理时间
c.put("max.poll.records", "50");               // 限制单批量,让处理可控
// 消费循环:处理完一批后 commitSync();失败不提交,等重试

broker 侧(topic 级,账务链路):replication.factor=3min.insync.replicas=2unclean.leader.election.enable=false

0.11 的幂等生产者(enable.idempotence=true)用于处理 producer 重试导致的重复。开启后由 broker 端去重,单个会话的单分区内不会因重试产生重复。producer 重启后去重状态丢失,且该机制只覆盖单会话、单分区。Kafka 0.11 也支持事务和精确一次语义,其外部系统写入限制见前文。当时团队依赖业务侧幂等(去重表),未使用幂等生产者和事务;对这部分的理解仅来自文档。

边界清单

  1. 分区内顺序。 同一分区内消息按写入顺序排列,消费者按顺序消费。跨分区没有全局顺序。
  2. 至少一次投递的配置前提。 acks=allmin.insync.replicas≥2、手动提交(先处理后提交)和消费幂等共同决定该语义。消息可能重复,重复消息由幂等处理。
  3. 跨分区顺序。 Kafka 不提供跨分区的全局顺序。要求全局顺序时,需要使用单分区,消费并发降为一。
  4. 业务侧恰好一次。 端到端 exactly-once 需要生产端幂等、事务,以及消费侧位移与处理原子化。Kafka 0.11 的事务不覆盖外部系统写入,多数场景使用至少一次投递和幂等处理。
  5. 积压消息的保留期限。 retention 到期后消息被删除,消费者长期落后会丢失消息。积压需要通过监控和扩容消费者处理,不在 Kafka 的保证范围内。
  6. 配置组合。 仅开启 acks=all、不设置 min.insync.replicas 时,ISR 收缩会降低确认条件;仅关闭 unclean、不设置 acks=all 时,leader 切换仍可能丢失消息;仅关闭自动提交、不配合幂等处理时,重复消息无法处理。可靠性取决于各环节配置,需要逐项核对。

业务侧处理重复投递的方法见消息语义:按 at-least-once 设计业务代码

理解有限的部分

当时只知道 controller 选举、分区状态机(分区的 online / offline / under-replicated 状态迁移)和 ISR 伸缩存在相应的内部流程,尚未理解其细节。controller 依赖 ZK 做选举,ZK 会话超时会触发 controller 切换,这部分留到ZooKeeper/etcd:小数据强一致的协调服务讲 ZK / etcd 时补。

参考资料


573 字 · 62 段落
ximing

Follow onGitHub

相关文章