1促销后夜:Lag 两百万与重复扣款短信同时炸开
某订单中台在促销高峰后同时收到两类告警:监控显示 order-events Topic 的消费者组
order-processor 积压超过 200 万条,Lag 持续攀升;客服工单激增,用户反馈
「同一笔订单收到两条扣款短信」。值班同学第一反应是「加消费者实例」——扩容后 Lag 短暂下降又迅速反弹。
并行排查重复消费时发现:部分实例在 rebalance 后从较早 offset 重新拉取,而消费端缺少幂等校验,
下游支付网关被重复调用。
进一步 --describe 才看到:6 个分区中 5 个 Lag 在 2 万以内,唯独 partition-3 达到 198 万。
日志里同一 offset 的 JSON 解析 NPE 每秒重试十余次——典型 poison pill 占满该分区消费线程。
重复扣款则集中在滚动发布窗口:K8s terminationGracePeriodSeconds=30,单批处理需 45 秒,
旧 Pod 完成支付 RPC 后尚未 commitSync 即被 SIGTERM,新 Pod 重放约 800 条,其中 23 条触发重复扣款。
这段复合场景暴露的不是「Kafka 不可靠」,而是团队把三类不同问题揉成一个口号:
Broker 层是否可能丢消息(会话、ISR、acks)、发送层是否因重试产生重复 Record(Producer 幂等)、
业务层是否因 at-least-once 与发布窗口产生重复处理(消费幂等 + rebalance 治理)。
积压要消「吞吐与阻塞」,重复消费要消「语义与窗口」——处置路径不同,却常在同一晚同时出现。
证据分层也要说在前面:副本、ISR、acks 行为属于与官方文档一致的机制共识;分区估算与阈值来自作者经验,
必须结合你们集群压测校准;案例均为复合场景,用于演示排查思路,不代表任何单一真实事故。
现场处置优先级(事中 35 分钟窗口)
故障窗口内建议严格按下列顺序动作,避免「拆东墙补西墙」:先确认集群是否还具备写入与副本健康, 再谈消费侧扩容;先看分区级 Lag 分布,再决定是加实例还是隔离单条;先止血用户可感知的重复副作用, 再做耗时的根因分析。顺序错了,最常见的后果是:一边加消费者一边打爆下游数据库,或一边查幂等一边 poison pill 继续堵住分区。
- 确认 Broker 与 ISR 健康(Under-Replicated Partitions 是否大于 0)。
- 定位 Lag 是否 skew:热点 key 还是单分区 poison pill。
- 检查消费组是否频繁 join/leave(rebalance 风暴)。
- 区分「整体产能不足」与「单条失败无限重试」。
- 重复客诉先止血(业务幂等开关 / 暂停非关键消费),再查 offset 与 Producer 配置。
| 现场现象 | 可能的架构缺口 | 本文章节 |
|---|---|---|
| Lag 全分区均匀上升 | 消费产能不足或生产突增 | 第 8 章 SOP |
| 单分区 Lag 极高 | 热点 key 或 poison pill | 第 6、8 章 |
| Under-Replicated 告警 | 副本 / ISR / Broker 资源 | 第 3、4 章 |
| 用户感知重复扣款 | 消费非幂等 + rebalance 窗口 | 第 5、6 章 |
| Broker 写入失败峰值 | acks / min.insync 配置不当 | 第 4 章 |
「扩容消费者就能消积压」是最常见的误判。分区数是消费并行度的硬上限;若 Topic 仅 6 分区且已有 6 个消费者,
再加实例只会 idle。先算 分区数 × 单分区吞吐上限,再决定加分区、改批处理还是隔离 poison pill。
复合场景值班时间线(为何「并行排查」)
复盘常见时间线大致如下:T+0 收到 Lag 告警;T+8 分钟扩容 3 个 Consumer 实例无效;
T+15 分钟 --describe 发现 partition-3 skew;T+22 分钟日志定位到同 offset 反复 NPE;
T+28 分钟将该消息写入 DLQ 并 skip;T+35 分钟 Lag 开始下降。重复扣款客诉在 T+40 分钟由支付团队开启
「商户订单号幂等」临时开关止血。整段窗口说明:积压与重复消费必须并行排查,
但处置手段不同——前者靠 DLQ/吞吐优化,后者靠幂等与 rebalance 治理。事后若只修 partition-3 而不更新
检查表中的 DLQ 与幂等项,下一轮大促仍会在同一薄弱点复发。
从现场反推架构缺口时,建议固定用「现象 → 机制 → 参数」三层提问:例如「Lag 单分区 skew」→ 「该分区消费线程是否被单条阻塞」→「ErrorHandler 是否无限重试、有无 DLQ」。避免直接跳到「加机器」 而跳过机制层。本文后半部分的决策树与 SOP,就是把这套提问顺序写成可交接的值班动作。
2高可靠三层栈:Broker 不丢、发送不重、业务可幂等
在进入 ISR 细节之前,先建立一张概念地图:Kafka「高可靠」不是单一开关,而是三层叠加。 任意一层缺失,都会在促销窗口以 Lag、客诉或写入失败的形式暴露出来。
后文按「总览 → ISR 与副本 → 可靠性参数 → 幂等 → 死信与 rebalance → 适合/不适合 → SOP → 决策矩阵」展开。 参数默认值与行为以 Kafka 3.x 官方文档为准;产能算式与阈值属于待压测校准的经验推断。
读后文时建议始终带着一张「责任归属表」:消息是否已进入 ISR 确认范围,是平台与 Topic 配置的责任; 网络抖动是否会产生重复 Record,是 Producer 客户端配置的责任;用户是否感知到重复扣款或重复通知, 是消费应用与下游幂等的责任。值班时最浪费时间的模式,是三方互相指认「Kafka 又出问题了」,却不先回答 「问题落在哪一层」。第 1 章的现场映射表与第 9 章的决策矩阵,分别对应「事中快速归层」与「事后选型收束」。
另需澄清一个常见措辞:口语里的「Kafka 保证 exactly-once」往往混用了三种含义——Broker 对已提交数据的持久性、 幂等 Producer 对重试的去重、以及事务或业务层对端到端恰好一次的拼装。本文刻意拆开讲述,是为了避免在架构评审里 用一个营销词覆盖三类完全不同的工程控制点。
3分区、副本与 ISR:Leader 切换时什么真正被保护
Topic 在逻辑上由多个分区组成,每个分区是一条只追加日志。为高可用配置
replication.factor 个副本并分布在不同 Broker:同一时刻仅 Leader 负责读写,
Follower 通过 Fetch 拉取并维持相同 offset 顺序。每个 Record 在分区内有单调递增的 offset,
这是消费进度与副本同步的共同基准。Follower 不是简单「复制文件」,而是持续从 Leader 拉取并保持相同顺序;
Leader 收到写入后先写本地 log,再按 acks 级别等待 ISR 确认。Producer / Consumer 默认与 Leader 交互
(Kafka 3.x 虽有 Follower Fetch 等进阶能力,生产默认仍以 Leader 为主);Leader 宕机时,
Controller(KRaft quorum)从 ISR(In-Sync Replicas) 中选举新 Leader。
若 ISR 为空,行为取决于 unclean.leader.election.enable:默认 false 时宁可分区不可用,
也不从非 ISR 副本选主,以避免数据回退。
ISR 抖动与运维信号
磁盘 IO 饱和、网络分区、Broker 长 GC 都会把 Follower 踢出 ISR。应持续监控
UnderReplicatedPartitions、OfflineReplicasCount,以及
IsrShrinksPerSec / IsrExpandsPerSec。频繁 shrink/expand 说明 Broker 层不稳,
Consumer 扩容无法根治。在 Grafana 中可为每个 Broker 建立 ISR 成员数面板:若某 Topic 分区 ISR size
长期为 1,等价于该分区没有热备,应触发工单而非仅记录 INFO。Follower 追平后会重新加入 ISR,过程对客户端透明;
但若磁盘长期打满,ISR 可能长期只剩 Leader——此时 min.insync.replicas=2 会让所有
acks=all 写入不可用,Producer 收到 NOT_ENOUGH_REPLICAS。这是「保护数据」与
「可用性」之间的刻意权衡,不应在生产环境靠调低 min.insync 来「临时恢复写入」。
Kafka 3.x 推荐 KRaft 管理元数据,Leader 选举路径更短,但选举原则仍是优先从 ISR 选。
若原 Leader 所在 Broker 永久下线,Controller 将 ISR 中 lag 最小的 Follower 提升为 Leader,客户端经
Metadata 更新自动发现。选举期间对应分区短暂不可写(通常秒级),客户端需配置合理的
retries 与 delivery.timeout.ms,避免瞬时 NotLeaderForPartition
被业务层误判为永久失败。KRaft quorum 大小须为奇数(3/5);从 ZooKeeper 迁移时分区 Leader 分布会变化,
建议低峰执行并观察 ActiveControllerCount 始终为 1。
分区数规划:并行度硬上限
| 考量维度 | 建议 | 反例 |
|---|---|---|
| 消费并行度 | 分区数 ≥ 峰值消费者实例数 | 12 实例配 6 分区,半数 idle |
| 生产吞吐 | 按单分区写入上限估算分区下限 | 高吞吐 Topic 仅 2 分区 |
| 顺序性 | 同业务键路由同分区 | 需要有序却用随机 UUID 作 key |
| 副本成本 | 分区数 × RF = 副本总数 | 千级分区未规划磁盘与句柄 |
复合场景估算:峰值 6000 TPS、单条 2KB,目标单分区写入不超过约 5MB/s 时,写入侧至少约 3 分区; 若消费 P99 80ms、仅 6 分区,理论消费上限远低于生产峰值——瓶颈在消费逻辑而非 Broker。 把分区从 6 扩到 12 只能抬高并行度上限,仍可能不够。该算式需用你们集群压测报告替换假设。
生产还应配置 broker.rack,避免 RF=3 却三副本同机架。扩容 Broker 后用
kafka-reassign-partitions.sh 均衡副本,防止新节点空载、老节点磁盘打满。
分区数一旦创建后只增不减(减少分区通常需重建 Topic),规划阶段应保守估算:
既要覆盖峰值消费者并行度,又要避免小集群上堆积过多分区导致文件句柄与控制器压力。经验上单 Broker
分区规模需结合硬件与版本评估,不可照搬「越多越好」。副本成本同样要进账:分区数乘以 RF 才是真实副本总数,
直接影响磁盘、跨机架带宽与恢复时间。
热点分区是 ISR 健康之外另一个高频坑:Producer 使用固定 key(如「默认店铺 ID」)会把流量打进少数分区, 表现与 poison pill 类似——Lag skew——但根因不同。前者要改 key 策略或拆 Topic,后者要 DLQ。 值班时若只看到 skew 就一律当 poison pill,可能在错误方向上浪费整晚。
4acks、min.insync 与幂等默认:可靠性参数边界
「Broker 级不丢」依赖一组互相咬合的参数,而不是单独把 acks 调到 all。
验证层目标:在测试环境故意 kill Leader、断网 Follower,确认 Producer 行为与监控告警符合预期。
# Broker / Topic 侧基线(RF=3)
min.insync.replicas=2
unclean.leader.election.enable=false
# Producer 侧(Java 客户端 3.x)
acks=all
enable.idempotence=true
retries=2147483647
max.in.flight.requests.per.connection=5
# 消费侧
enable.auto.commit=false
isolation.level=read_committed # 配合事务 Producer
acks=all(或 -1)表示 Leader 等待所有 ISR 副本确认后才向 Producer 返回成功;
当 ISR 成员数低于 min.insync.replicas 时,写入会被拒绝。该行为与官方 Replication / Producer 文档一致。
验证层怎么做,而不是只对配置文件
参数表写进 Wiki 并不等于可靠。建议在预发或演练窗口固定做三类注入:随机 kill 当前分区 Leader,
确认 acks=all + min.insync.replicas=2 下已确认消息不丢、客户端可在重试后恢复;
人为让一个 Follower 长时间断网,观察 ISR 收缩、URP 告警与(在 ISR 不足时)写入失败是否符合预期;
对事务 Producer 场景再验证 isolation.level=read_committed 的消费者不会读到未提交消息。
演练记录应留下「期望行为 / 实际行为 / 偏差」三列,避免只截一张配置截图作为验收证据。
若演练中出现「配置正确但客户端仍报错」,优先查是否混用了旧版客户端默认值、是否有旁路 Producer
(脚本、临时任务、数据修复工具)绕过了统一配置中心——这类旁路是生产里最容易漏掉的可靠性缺口。
机架感知方面,Kafka 支持 broker.rack 与机架感知副本选择,使副本分布在不同机架或可用区。
若 RF=3 但三副本同 rack,单可用区故障会同时失去写入能力与数据冗余。创建 Topic 时可手动
--replica-assignment,日常更依赖自动分配;扩容 Broker 后务必 reassign,避免「新节点空载、
老节点磁盘 90%」的隐性单点。
5Producer 幂等与消费端幂等:各自解决哪一类「重复」
「重复」至少有两种:Broker 上出现两条相同业务内容的 Record(常由发送重试引起);以及同一条 Record 被业务逻辑处理了两次(崩溃未提交 offset、rebalance、下游超时重试)。前者靠 Producer 幂等(及事务), 后者必须靠消费端设计——开启前者绝不等于后者可省略。
| 消费端模式 | 实现要点 | 适用场景 |
|---|---|---|
| 天然幂等 | 下游可重复(覆盖写固定主键) | 状态同步、缓存刷新 |
| 业务唯一键 | DB 唯一索引 + ON CONFLICT | 订单、支付流水 |
| 消费日志表 | (group,topic,partition,offset) 唯一后执行业务 | 强一致账务 |
| Redis SETNX | msgId / orderId,TTL 覆盖重试窗口 | 高 QPS 通知 |
| Outbox + CDC | 本地事务写业务与 outbox,再发 Kafka | DB 与消息一致性 |
开启 Producer 幂等不等于消费端可以不防重。Producer 幂等解决发送重试导致的重复 Record; Consumer rebalance、下游超时重试仍会产生重复处理。同理,DLQ 只隔离 poison pill,不替代消费端幂等。
消费日志表与 Redis 去重如何选型
消费日志表是强一致场景最常用的模式:表含 (consumer_group, topic, partition, offset) 联合唯一键。
流程为:开 DB 事务 → INSERT 消费记录(冲突则说明已处理)→ 执行业务 → COMMIT → 再
commitSync Kafka offset。这样即使 commit offset 失败,重放时 INSERT 冲突会拦截重复业务。
表会随消费增长,需按 TTL 归档或按月分表。与 Outbox 对比:Outbox 侧重「本地事务发消息」,消费日志表侧重
「消费侧去重」;账务系统二者常组合使用。
Redis SETNX 方案须设置 TTL 大于 Consumer 最大重试窗口(常见建议按天级,例如覆盖 7 天回放窗口),
key 可设计为 dedup:{topic}:{group}:{partition}:{offset} 或业务 orderId。
前者严格按位点去重,后者按业务语义去重;支付场景更推荐业务键,因为 rebalance 后 offset 变了但订单号不变。
复合场景中重复扣款的直接原因往往是「调用支付 API → 提交 offset」顺序:API 成功但进程在 commit 前退出,
重启后重放。修正为「先写 dedup → 调 API → commit」,才能把窗口关掉。
事务 Producer 通过 transactional.id 在 Broker 注册会话,支持
sendOffsetsToTransaction 将消费 offset 与产出消息原子提交,适用于
consume-transform-produce。代价是吞吐下降与事务状态开销;订单主链路多数情况下
「幂等 Producer + 消费端去重」已足够,跨 Topic 强一致再上事务。
6死信队列、重试分层与 rebalance 窗口治理
消费失败若无限重试同一条消息,会阻塞该分区后续消息(head-of-line blocking),表现为 「整体 Lag 上升但错误日志只有一条」。必须区分可重试异常(下游 503)与不可重试异常(序列化失败、校验失败): 后者应尽快进 DLQ,避免耗尽重试预算。框架上至少三层:瞬时错误做有限次本地重试加指数退避; 可修复错误重试 N 次后进死信;永远失败的 poison pill 跳过或进 DLQ,人工修复后再回放。 缺少分层时,ErrorHandler 会把「一条坏消息」放大成「整组不可用」。 这正是复合场景里 partition-3 拖垮整组 Lag 的典型机制,必须优先隔离。
DeadLetterPublishingRecoverer 可按原 partition 映射,便于分区级追踪。
Rebalance 窗口:重复消费的放大器
Kafka 3.x 默认含 CooperativeStickyAssignor,可减少全量停顿。仍建议:
max.poll.interval.ms 覆盖单批 P99;session.timeout.ms 与
heartbeat.interval.ms 约 3:1;静态成员 group.instance.id 降低滚动发布抖动。
复合场景复盘:将 max.poll.interval 调到 15 分钟、max.poll.records 降到 100,单批处理从约
45 秒降到约 12 秒,并配置 group.instance.id=${HOSTNAME} 与 K8s preStop,滚动发布时
rebalance 次数显著下降,重复扣款未再出现——属生产经验区间,须按你们批处理耗时与发布策略校准,
不可直接照搬数值。
DLQ Topic 命名常见约定为 {原Topic}.DLQ 或 {原Topic}.dead-letter,RF 与主 Topic
一致;权限上仅运维与开发角色可消费,避免误触发业务逻辑。若使用 Spring Kafka 的
DeadLetterPublishingRecoverer,默认可把原 partition 映射到 DLQ 同号分区,便于分区级追踪;
也可集中路由到 partition 0 供人工消费——前者利于追踪,后者利于运维但可能成热点。回放工具从 DLQ 读取后
写入主 Topic 或专用 replay Topic,消费端必须再次走幂等校验。禁止把「跳过 poison pill」理解成删除主 Topic
数据或未审批的 --reset-offsets --to-latest。
Consumer 跑在 K8s 时,务必把 terminationGracePeriodSeconds 与最长批处理时间对齐,
并在 preStop 中给足 commit / 退组时间。否则「支付成功未 commit」会在每一次滚动发布复现。
该项应与 K8s 生产稳定性检查表联动巡检:探针、优雅退出与 Kafka 的 max.poll / session 超时是同一条链路,
拆成两个团队各自配置时最容易出现「两边都觉得自己合理、合在一起必炸」的窗口。
7什么场景该押注 Kafka 高可靠栈,什么时候别硬上
高可靠参数组合会牺牲延迟与部分可用性(ISR 不足时拒绝写入)。不是所有 Topic 都该套「金融级」基线, 也不是所有「异步解耦」都适合用 Kafka 充当事务总线。
更适合
- 订单、支付、库存等不可丢且可异步的领域事件
- 需要分区内有序、按 key 扩展并行度的流水线
- 多消费者组独立进度(同一 Topic 服务多个下游)
- 已具备监控、DLQ、幂等与 on-call SOP 的平台团队
- 吞吐波动大、需要削峰填谷的写入路径
不适合 / 需谨慎
- 要求同步强一致、毫秒级请求-响应的在线交易主路径(应 RPC/本地事务)
- 允许丢失的海量日志却套 acks=all + 事务,成本无收益
- 团队无力维护 DLQ 回放与消费幂等,却声称「exactly-once」
- 用 Kafka 当数据库:超大消息、随机读改、长期点查
- 分区数与机架拓扑未规划就上多活跨机房(先单集群可靠)
若业务真正需要的是「跨 Topic 原子写出 + 消费 offset 同事务提交」,再评估 Kafka 事务; 多数订单主链路用「幂等 Producer + 消费端业务键去重」性价比更高。事务会降低吞吐并增加 Broker 状态开销—— 这是机制代价,不是配置疏忽。
与「消息中间件万能论」划清边界
把同步下单接口改成「写 Kafka 再异步落库」可以削峰,但也会把用户可见的失败模式从「接口报错」变成 「下单成功短信延迟 / 状态短暂不一致」。若产品与客服体系无法接受最终一致窗口,就不该为了「架构好看」硬上。 反之,对已经异步化的领域事件,若仍坚持用数据库轮询或点对点队列且缺少分区扩展能力,则 Kafka 的多消费者组 与水平扩展往往更合适。适合与不适合的判断,最终要回到一致性窗口、运维成熟度、峰值形状 三件事,而不是回到「业界是不是都在用 Kafka」。
团队能力也是硬约束:没有 DLQ 回放、没有消费幂等测试、没有 URP/Lag 值班手册时,把
acks 调到 all 只会制造「配置看起来很稳、故障时无人会修」的假象。此时更务实的路径是先补齐
第 8 章检查表中的监控与 DLQ 两项,再谈事务与多活。
8积压排查决策树、可靠性检查表与标准 SOP
框架层要把经验固化为可执行步骤:积压归因四类(产能不足、生产突增、poison pill、rebalance 风暴), 再对照扩容与优化取舍。混合根因很常见——分别估算「生产恢复常态后消化时间」与「当前消费 TPS 上限」, 取更悲观者作为 SLA 承诺。SOP 与检查表的分工同样要说清:SOP 管事中十五分钟到两小时的动作顺序, 检查表管事前季度巡检与上线门禁;缺少任何一张,都会在大促夜把经验重新变成口头相传。
扩容 vs 优化取舍
消费者扩容是「买时间」,逻辑优化是「治本」。架构评审应要求提供产能算式:消费 TPS 上限粗算约为
分区数 × (1000ms / 单条处理毫秒) × 批处理条数(未计 rebalance 与 GC)。若算式结果低于业务峰值,
扩容只能延缓积压。实际排障中常出现混合根因:例如生产突增两倍叠加消费逻辑变慢,Lag 呈指数而非线性——
应分别估算「生产恢复常态后多久能消化」与「当前消费上限」,取更悲观者作为对外 SLA。
| 条件 | 扩容消费者 | 优化消费逻辑 |
|---|---|---|
| 消费者数 < 分区数 | 有效,优先加实例 | 并行进行,长期仍需优化 |
| 消费者数 = 分区数 | 无效,须先加分区 | 降单条耗时、批量写 |
| 下游 DB 瓶颈 | 可能加剧压力 | 批量、异步、读从库 |
| 单条解析失败循环 | 无效 | DLQ 隔离 poison pill |
与监控指标联动
除 Lag 外,Dashboard 建议固定展示:Broker 侧 MessagesInPerSec 识别生产突增;
Fetch 相关请求耗时识别消费拉取是否被拖慢;客户端 records-lag-max 与
bytes-consumed-rate;应用层单条处理 P99 与 DLQ 写入速率。当 MessagesInPerSec 平稳而 Lag
上升,几乎可排除 Producer 突增,直接进入 poison 或下游瓶颈分支。反之若生产速率与 Lag 同步上升,
先与业务确认大促、批量补数、CDC 全量等已知事件,再决定是否临时限流 Producer。
排障命令(Kafka 3.x)
bin/kafka-consumer-groups.sh --bootstrap-server kafka1:9092 \
--describe --group order-processor
bin/kafka-topics.sh --bootstrap-server kafka1:9092 \
--describe --topic order-events
# 重置 offset 到最新(慎用,会丢未消费消息,须审批)
bin/kafka-consumer-groups.sh --bootstrap-server kafka1:9092 \
--group order-processor --topic order-events \
--reset-offsets --to-latest --execute
可靠性检查表(季度 / 重大变更巡检)
评分:符合 1 分、不符合 0 分、不适用不计;核心链路 Topic 得分低于 80% 禁止上线新 Consumer。 新建 Topic 另做「可靠性三板斧」:kill Leader 不丢数;进程 kill 不产生不可接受重复;人工 poison pill 5 分钟内进 DLQ 且主链路恢复。
| 检查项 | 标准 | 责任 |
|---|---|---|
| RF | 生产 Topic RF≥3,跨机架/可用区 | 平台 |
| min.insync.replicas | RF=3 时 ≥2 | 平台 |
| unclean.leader.election | 集群级 false | 平台 |
| Producer acks / 幂等 | 核心链路 acks=all + enable.idempotence | 开发 |
| 消费提交与幂等 | 手动 commit;业务唯一键或消费日志表 | 开发 |
| 死信队列 | 失败 N 次进 DLQ,retention≥14 天,有回放 SOP | 开发 |
| 分区与监控 | 分区≥峰值消费者;Lag/URP/延迟有 on-call | 架构 / SRE |
| Rebalance / Offset | max.poll 覆盖 P99;reset-offsets 需审批 | 开发 / SRE |
积压排查 SOP(八阶段)
| 阶段 | 动作 | 判读要点 |
|---|---|---|
| 1. 现象确认 | 记录 Lag、Topic/Group、起止时间 | 持续上升 vs 瞬时尖峰 |
| 2. Broker 健康 | 查 URP、Offline、Controller | URP>0 先修副本 |
| 3. 分区 skew | 对比各 partition Lag | 单分区占 80%+ 查 key/poison |
| 4. 消费组状态 | 成员数、rebalance 频率 | 频繁 JoinGroup 查慢 poll |
| 5. 根因定位 | 对照决策树四类 | CPU 低 Lag 高常见 poison |
| 6. 处置 | 扩容 / DLQ / 限流 / 批处理 | 禁止未审批 to-latest |
| 7. 验证 | Lag 下降并稳定 30min+ | 反弹则回步骤 5 |
| 8. 预防 | 更新检查表与演练记录 | 大促前下调阈值并演练 DLQ |
步骤 6 中「跳过 poison pill」仅指将该 offset 消息移入 DLQ 并提交下一 offset,或修复后回放; 不是删除主 Topic 数据。任何 offset 跳跃须第二人复核,变更单记录 partition、offset、操作人与原因。 重大故障结束后 48 小时内做 SOP 符合性回顾,标记被跳过的步骤。
复合场景完整复盘:partition-3 的 poison pill 入 DLQ 后,其余分区本身无积压,Group Lag 迅速回落; 团队随后做 schema 校验直进 DLQ、分区 6→12 并扩消费者、支付侧商户订单号状态机幂等。二次演练中 Lag 峰值可控且可在约 20 分钟内自行消化。说明SOP 管当下,检查表管下次不再—— 积压与重复看似独立,共享根因是消费框架缺少 DLQ 与幂等基线。
9决策矩阵:继续加固、降级语义还是换路径
当评审卡在「要不要把所有 Topic 都调到金融级」或「要不要上事务 / 多活」时,用下面的决策矩阵收束讨论。 列含义:继续(在现有 Kafka 栈上加固)、迁移(换语义或换组件)、选型(新建系统时的默认选择)。
| 业务信号 | 继续(加固 Kafka) | 迁移 / 降级 | 选型建议 |
|---|---|---|---|
| 订单/支付事件,可异步 | RF=3、acks=all、幂等、DLQ、消费幂等 | 仅当必须同步强一致时改本地事务+Outbox | Kafka 高可靠基线为默认 |
| 可丢失的观测日志 | 可保留 Kafka 但降 acks / 缩短 retention | 可迁对象存储或轻量采集链路 | 勿套金融级参数 |
| 跨 Topic 原子写出 | 评估 transactional.id 与吞吐代价 | Outbox/Saga 可能更清晰 | 先证明需要再上事务 |
| Lag 反复、分区已满配 | 优化批处理 + DLQ + 有规划加分区 | 下游改异步化或拆 Topic | 禁止只加无分区的消费者 |
| 重复客诉与发布强相关 | 静态成员、grace、消费幂等 | 发布策略改为蓝绿+排水 | K8s 与 Kafka 参数联动 |
| 多活 / 跨机房 | 先单集群 ISR/机架达标 | 再评估 MirrorMaker 2 等复制 | 可靠单集群优于半吊子双活 |
三句话收束
- Broker 级可靠:RF、ISR、acks=all、min.insync.replicas,拒绝 unclean 选举。
- 发送级可靠:Producer 幂等 + 合理 key;全局恰好一次另议事务或业务键。
- 业务级可靠:手动 commit、消费幂等、DLQ;积压先诊断再扩容,分区数是并行度硬上限。
回到开篇的风格比例:问题现场层告诉我们积压与重复消费常同时出现但处置路径不同;体系架构层给出 ISR、幂等、DLQ 的约束边界;框架层用决策矩阵、检查表与 SOP 把经验固化为可执行步骤。Kafka 高可靠 不是某一个参数的开关能力,而是三层叠加后,在促销窗口仍然站得住的工程结果。
相邻知识地图
- 系列 T01 Hadoop 生态中 Kafka 作为数据管道角色;数仓分层实践中 CDC 入 ODS 的可靠性边界可与本文对照。
- T20 Redis 高阶可对比「缓存一致性 vs 消息最终一致」:写路径先落 Kafka、异步更新 DB 并失效缓存是常见组合。
- 架构师系列中的 Kafka 容灾 / 万亿消息篇可承接 MirrorMaker 2、跨机房 RPO/RTO(超出本文单集群范围)。
- Consumer 跑在 K8s 时,与生产稳定性配置中的探针、优雅退出、资源配额章节交叉检查。
- 进阶方向:Kafka Connect 与 CDC 链路、Tiered Storage 对重放的影响、多活下的 offset 映射——同构问题仍是副本、幂等与 DLQ。
分享讨论题
- 核心 Topic 的 RF 与 min.insync.replicas 实际是多少?是否仍有 acks=1 的订单 Producer?
- 消费端幂等用的是哪种模式?最近一次 DLQ 回放演练在什么时候?
- Lag 告警时谁有权执行
--reset-offsets?审批是否书面化? - 最高 QPS 的三个消费组,能否当场写出「分区数、消费者数、单条 P99、理论 TPS 上限」四元组?
- 上一次滚动发布是否与重复消费客诉时间线重合?grace 与 max.poll 是否联合评审过?
建议在完成本文后,结合团队现有 Topic 清单做一次检查表打分,并对最高流量的消费组演练积压 SOP: 人为构造一条不可解析消息,验证是否在约定时间内进入 DLQ 且主链路恢复;再在低峰做一次受控的 Leader 故障演练,核对监控与 Producer 重试行为。能把「文档里的可靠」变成「演练过的可靠」, 这篇分享才算落到架构治理,而不是又一次参数抄写。
最后再补一条可执行习惯:每次变更 RF、min.insync.replicas、acks 或 enable.idempotence 时, 同步更新 Topic 级变更单与告警阈值,并在复盘纪要里写明「预期一致性窗口」——这样下次复合故障 到场时,不必从零推断参数意图。