一句话回答
先确认堆积范围、增长速度和业务影响并停止流量放大,再区分生产突增、消费者异常、处理变慢、分区倾斜或下游瓶颈;短期通过恢复消费者、限流和增加有效并行度止损,长期通过容量规划、批量处理、背压和降级机制防止复发。
面试考察点
- 是否知道 Lag 是结果,不是根因。
- 能否区分生产速率、消费速率和预计清空时间。
- 是否理解消费者数量超过分区数不会继续提高并行度。
- 能否处理消息即将超过保留时间的紧急情况。
- 是否具备限流、扩容、数据转移和恢复后的完整方案。
什么叫消息堆积
对某个消费者组,最新日志位点与已提交消费位点之间的差值就是 Lag。Lag 偶尔上升不一定异常,关键看它是否持续增长以及业务能接受的处理延迟。
积压增长速度 = 生产速率 - 消费速率
预计清空时间 = 当前积压量 / 恢复后的净消费速率
如果生产每秒 10 万条、消费每秒 8 万条,Lag 每秒增加 2 万。即使消费者仍在工作,系统也处于持续失血状态。
第一步:确认影响范围
收到告警后先回答:
- 哪个集群、Topic、消费者组和分区堆积?
- Lag 是均匀增长还是集中在少数分区?
- 消费者实例是否存活,活跃成员数量是否变化?
- 生产速率是否突增,消费速率是否下降?
- 最老未消费消息已经延迟多久?
- Topic 保留时间还剩多少安全窗口?
- 哪些业务结果、通知或数据同步受到影响?
只看消费者组总 Lag 会掩盖分区倾斜。必须下钻到 Partition。
第二步:先止损
1. 停止重试风暴
下游超时后立即无限重试,会让同一批消息反复占用线程。应使用有限重试、指数退避,并把持续失败消息送入重试 Topic 或死信队列。
2. 限制上游非关键流量
若系统已经过载,可以暂停低优先级事件、降低采样率或让生产端缓冲。不能控制上游时,也要避免消费者因过载频繁崩溃和再均衡。
3. 隔离毒消息
某条格式错误或业务异常消息可能反复失败,阻塞整个分区。记录原消息、异常和重试次数,超过阈值后转入死信队列,让主消费链路继续推进。
资金等关键消息不能直接跳过,必须建立告警和人工修复流程。
常见根因一:消费者进程异常
检查消费者是否频繁重启、OOM、Full GC、线程池耗尽或发生未捕获异常。若实例不健康,单纯扩容只会增加更多失败实例。
同时检查再均衡频率。消费者处理超过 max.poll.interval.ms,会被移出组并反复再均衡,表现为实例都存活但实际吞吐很低。
常见根因二:下游依赖变慢
消费者通常要写 MySQL、Redis、ES 或调用接口。下游连接池耗尽、慢 SQL、锁等待和限流都会把消费线程阻塞。
应对方法:
- 先恢复或降级下游服务。
- 批量写数据库,减少网络往返。
- 缩短事务并补充必要索引。
- 设置合理超时,避免线程永久等待。
- 对非关键下游采用异步或暂存。
如果下游每秒只能处理 1 万条,把 Kafka 消费者扩到每秒拉取 10 万条只会把压力转移并拖垮下游。
常见根因三:生产流量突增
活动、爬虫、批量回放或上游 Bug 可能让生产量突然翻倍。应确认是正常峰值还是异常流量,并判断容量是否足以在业务时限内追平。
正常峰值需要预先扩容和压测;异常生产需要限流、修复生产者并识别已经产生的无效消息。
常见根因四:分区倾斜
总消费能力充足,但某个热点 Key 把大量消息写入一个分区。该分区只能由组内一个消费者处理,其他消费者空闲。
常见原因:
- 分区 Key 分布不均。
- 大客户或超级租户占据单个 Key。
- 空 Key 被统一路由。
- 自定义分区器存在缺陷。
修改分区策略会影响同 Key 有序性。可以拆分热点 Key、为超级租户使用独立 Topic,或在业务允许时增加更细粒度的并行键。
常见根因五:分区数不足
消费者组内一个分区同时只分配给一个消费者。若有 12 个分区,增加到 30 个消费者,仍只有 12 个消费者工作。
增加分区可以提高并行上限,但要注意:
- 旧 Key 到分区的映射可能变化。
- 全局或 Key 有序语义需要重新评估。
- Broker 元数据、文件句柄和复制开销增加。
- 已经存在的积压不会自动均匀迁移到新分区。
因此临时增加分区并不一定能快速清掉旧积压。
常见根因六:单条消息处理太慢
消费者可能逐条查询、逐条写入,形成 N+1。可以在不破坏顺序和事务边界的前提下批量处理:
单条写数据库 1000 次
↓
按分区聚合后批量写 10 次
还可以缓存重复读取的维度数据、减少日志同步输出、预编译表达式和优化序列化。
如何临时扩容
横向增加消费者
当活跃消费者数少于分区数,且瓶颈在消费者 CPU 或本地处理时,增加实例通常有效。扩容后观察实际分区分配、消费速率和下游压力。
提高单实例并行度
poll 线程可把不同分区任务交给工作线程,但必须维护分区内顺序和位点水位。不能任务一提交线程池就立刻提交位点。
建立临时快速消费组
极端积压时,可以将消息转存到临时 Topic,使用更多分区和专用消费者快速处理。但这属于高风险方案,需要确保不丢失、不乱序、可对账,并限制临时链路重复处理。
消息快过期了怎么办
当最老消息接近 Topic 保留期限,优先扩大安全窗口:
- 评估并延长 Topic 保留时间。
- 确认 Broker 磁盘容量足够。
- 将未消费数据复制或导出到临时存储。
- 暂停可能加速磁盘耗尽的非关键流量。
- 建立消息 ID 清单和对账任务。
直接无限延长保留时间可能填满磁盘,引发更大故障。必须同时计算积压字节量、增长率和剩余空间。
恢复后如何安全追平
恢复消费能力后,不要立即解除所有限流。持续观察:
- Lag 是否稳定下降。
- 消费错误率和重试量。
- MySQL、Redis、ES 等下游容量。
- 消费者 GC、CPU、线程池和连接池。
- Broker 网络、磁盘和请求延迟。
追平期间可能大量触发过期业务,例如延迟数小时的通知、优惠券或风控事件。消费者要判断消息时效,决定执行、丢弃还是补偿。
预计清空时间怎么计算
假设当前积压 3.6 亿条,生产速率 2 万条/秒,扩容后消费速率 5 万条/秒:
净清理速度 = 5 万 - 2 万 = 3 万条/秒
预计清空时间 = 3.6 亿 / 3 万 = 12000 秒 ≈ 3.3 小时
实际还要考虑流量波动、失败重试和下游限流。清空目标必须使用净消费速率,不能直接用总消费速率。
长期治理方案
- 为每个消费者组设置 Lag、最老消息年龄和消费速率告警。
- 压测峰值生产与消费能力,保留安全余量。
- 为下游依赖设置超时、隔离、限流和降级。
- 建立重试 Topic、死信队列和人工补偿平台。
- 监控分区倾斜和消费者再均衡。
- 规范批量回放和人工重置位点操作。
- 为消息设置业务时效和过期处理策略。
- 定期演练消费者全停、下游故障和大规模积压。
常见误区
- 看到 Lag 就无限增加消费者,忽略分区数和下游容量。
- 为清积压直接提交最新位点,相当于主动丢弃所有历史消息。
- 无限重试毒消息,导致一个分区永久阻塞。
- 只看总 Lag,不检查单分区热点。
- 追平后立即缩容,没有观察重试和业务补偿。
- 延长保留时间却不计算磁盘剩余空间。
核心考点清单
- Lag 持续增长说明生产速率长期高于有效消费速率。
- 排查要下钻到消费者组和分区,并区分生产突增与消费下降。
- 消费者扩容上限受分区数限制,也受下游容量限制。
- 毒消息应隔离,关键消息必须进入可追踪的补偿流程。
- 清理时间使用“消费速率减生产速率”的净速度计算。
- 消息接近过期时,要同时处理保留时间、磁盘和数据备份。
高频追问与参考回答
追问 1:消息堆积后直接增加消费者有用吗?
只有消费者数少于分区数、瓶颈确实在消费者且下游还有容量时才有效。消费者超过分区数会空闲,下游已饱和时扩容反而加剧故障。
追问 2:某一个分区积压特别严重怎么办?
检查分区 Key 和热点业务。短期可为热点业务扩容单个处理链路或转移数据,长期要拆分热点 Key、调整分区策略或隔离超级租户,同时评估顺序语义。
追问 3:能不能重置位点到最新位置快速恢复?
这会跳过全部积压,相当于丢消息。只有业务明确允许丢弃、数据已备份且完成审批对账时才能使用,关键消息不能这样处理。
追问 4:如何判断多久可以清完?
使用当前积压量除以净清理速度,即扩容后消费速率减去同时发生的生产速率,并为重试和波动预留余量。
追问 5:消费者没有报错为什么还会积压?
可能是处理速度稳定但低于生产速度、某分区热点、下游响应变慢、批次过小或频繁再均衡。没有异常日志不代表吞吐满足需求。
追问 6:积压恢复后要做什么?
核对消息与业务结果、处理死信和过期事件,逐步解除限流,复盘容量和告警缺口,并通过压测与故障演练验证修复。
机制全景图
「Kafka 消息堆积时应该怎么处理?」的实现链路如下,节点可与后面的源码和运行证据逐一对应。
flowchart LR
A["发现 Consumer Lag 上升"]
A --> B["判断生产增速或消费降速"]
B --> C["定位分区/实例瓶颈"]
C --> D["止损扩容或降级"]
D --> E["估算追赶时间并核对结果"]
源码与实现定位
| 入口 | 阅读重点 |
|---|---|
| records-lag-max | 分区最大积压 |
| consumer-fetch-manager-metrics | 拉取与消费速率 |
源码或系统表应按上表顺序追踪:先确认入口实际走到哪条路径,再用运行时数据验证,而不是仅凭类名或配置推测。
参数配置与可复现实验
catch_up_seconds = lag / (consume_rate - produce_rate)
制造热分区、毒消息和慢数据库,分别验证追赶公式。
验证步骤与预期结果
1. 固定输入和基线
先在没有故障注入的环境执行上述配置,固定数据规模、并发度、运行时版本和预热时间。以「净消费速率」为主基线,记录值应满足「记录分区基线」;同时保存 各分区 Lag、生产/消费净速率,使后续变化能够回到同一时间轴比较。
2. 从实现入口确认路径
在「records-lag-max」确认请求确实进入「分区最大积压」对应的实现,再沿「consumer-fetch-manager-metrics」观察「拉取与消费速率」。如果入口路径都未命中,就不应继续调整下游参数,而应先检查调用条件、版本或路由是否与假设一致。
3. 注入本文特有的失败模式
优先复现「总 Lag 掩盖单分区热点」,并把单一变量逐级放大,直到「净消费速率」越过「超过基线 2 倍」。随后再分别验证「无限重试毒消息阻塞分区」和「扩容消费者压垮下游」,三类故障分开执行,避免多个变量同时变化而无法归因。
4. 执行止损和根因修复
第一轮只应用「先恢复正净速率」,确认它能控制影响范围;第二轮应用「隔离毒消息」,验证核心链路恢复;最后落实「扩容不超过分区/下游容量」,消除同类问题再次出现的条件。每一步都保留变更前后数据,不用“感觉变快了”替代测量。
5. 通过退出条件
实验只有同时满足三项才算通过:「净消费速率」回到「记录分区基线」、「P99 延迟」回到「小于业务预算」、「端到端差异」回到「0」,并且业务结果差异为零。若性能恢复但结果不一致,仍应视为失败;若指标恢复后很快再次越线,则说明只完成了临时止损,没有消除根因。
量化基线
| 指标 | 样例基线/口径 | 风险线 | 结论 |
|---|---|---|---|
| 净消费速率 | 记录分区基线 | 超过基线 2 倍 | 无法追赶 |
| P99 延迟 | 小于业务预算 | 突破预算 | 无法追赶 |
| 端到端差异 | 0 | 任意非零 | 停止并对账 |
这些数值是实验口径或示例告警线,不是可复制到所有系统的固定答案;上线阈值应由本系统稳态、峰值和故障演练共同确定。
事故复盘:订单事件积压数小时仍无法追平
消费者扩容后数据库连接池已饱和,更多实例只增加等待,净消费能力未提升。优化批量写、提高分区并按数据库可承载并发设置消费者后,才形成正的追赶速率。
| 失败模式 | 首要证据 | 第一处置动作 |
|---|---|---|
| 总 Lag 掩盖单分区热点 | 各分区 Lag | 先恢复正净速率 |
| 无限重试毒消息阻塞分区 | 生产/消费净速率 | 隔离毒消息 |
| 扩容消费者压垮下游 | 单批处理时长 | 扩容不超过分区/下游容量 |
发布与回滚检查点
- 发布前:确认「records-lag-max」对应实现和上述配置在目标版本仍然有效,并保存「净消费速率」基线。
- 灰度中:同时观察 各分区 Lag、生产/消费净速率、单批处理时长;任一指标越过表中风险线,就停止继续扩量。
- 回滚时:先执行「先恢复正净速率」控制影响,再回退代码或参数;涉及持久状态时必须额外核对结果差异。
- 发布后:至少覆盖一个完整峰值周期,确认「总 Lag 掩盖单分区热点」没有再次出现,才关闭变更观察窗口。
设计边界与工程取舍
积压恢复必须同时保证正确性;跳过、改 offset 或扩大批次都要明确重复、丢失和乱序的补偿方式。