Kafka 消息积压处理
这是 Kafka 运维面试中非常高频的一道题。一个回答不能只说“增加消费者”,需要体现定位原因 → 临时止血 → 根因解决 → 防止复发的完整思路。
1. Kafka 消息积压怎么处理?
2. 先说明什么叫消息积压
Kafka 消息积压(Consumer Lag)指:
Producer 生产的消息速度 > Consumer 消费消息速度,导致 Topic 中未消费消息越来越多。
Kafka 中通常通过 Consumer Lag 衡量:
Lag = Log End Offset - Consumer Offset
例如:
Topic partition:
生产位置:
100000
消费者消费到:
80000
Lag:
20000
说明还有 2 万条消息未消费。
3. 第一步:确认是否真的积压
4. 查看 Consumer Group 状态
kafka-consumer-groups.sh \
--bootstrap-server kafka01:9092 \
--describe \
--group order-consumer
输出:
TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG
order 0 80000 100000 20000
order 1 90000 100000 10000
关注:
-
CURRENT-OFFSET
- 消费者当前消费位置
-
LOG-END-OFFSET
- Kafka 最新消息位置
-
LAG
- 积压数量
5. 分析积压原因
Kafka 消息积压通常有 5 类原因。
6. 消费者消费能力不足(最常见)
7. 表现:
生产速度:
10000 msg/s
消费速度:
5000 msg/s
Lag持续增长。
8. 排查:
查看消费者数量:
kafka-consumer-groups.sh \
--describe \
--group xxx
例如:
Topic partition:
partition 0
partition 1
partition 2
partition 3
但是:
consumer:
2个
说明:
消费并发不足。
9. 解决:
增加 Consumer 实例。
例如:
原来:
Topic:
4 partitions
Consumer:
2个
调整:
Consumer:
4个
提高:
并行消费能力。
注意:
消费者数量不能超过 partition 数量。
例如:
Topic:
3 partitions
Consumer:
10个
实际:
只有3个有效消费。
10. 消费者处理消息太慢
例如:
业务:
收到订单消息
↓
调用库存接口
↓
调用支付接口
↓
写数据库
单条消息耗时:
500ms
吞吐:
1秒只能消费2条
11. 排查:
查看:
-
consumer CPU
-
GC
-
DB慢查询
-
外部接口
例如:
应用日志:
process message cost 2000ms
12. 优化:
12.1. 批量消费
调整:
max.poll.records=500
一次拉取更多消息。
12.2. 异步处理
同步:
Kafka
↓
业务处理
↓
commit offset
异步:
Kafka
↓
线程池
↓
业务处理
提高吞吐。
12.3. 优化业务逻辑
例如:
不要:
每条消息:
insert mysql
改:
批量:
insert 1000条
13. Kafka 分区设计不合理
14. 问题:
Topic:
partition=3
但是:
生产:
100万消息/s
消费:
无法扩展。
15. 解决:
增加 partition:
kafka-topics.sh \
--alter \
--topic order \
--partitions 12
变成:
12 partitions
消费者:
可以扩展到12个。
注意:
partition 增加后:
-
不影响已有数据
-
新消息重新分配
-
key 顺序可能变化
16. 消费者异常停止
例如:
服务挂了:
Consumer Down
↓
没有消费
↓
Lag增长
排查:
查看消费者:
kafka-consumer-groups.sh \
--describe
发现:
Consumer ID:
empty
说明:
没有消费者。
检查:
应用:
systemctl status app
或者 Kubernetes:
kubectl get pod
查看:
kubectl logs pod-name
17. Kafka Broker性能不足
如果:
消费者正常
但是:
读取速度慢。
检查 Broker。
18. 查看:
CPU:
top
磁盘:
iostat -x 1
重点:
await
util
如果:
util 100%
说明磁盘瓶颈。
网络:
sar -n DEV 1
解决:
-
增加 Broker
-
增加磁盘
-
SSD替换机械盘
-
增加副本合理分布
19. 紧急情况下如何快速处理积压?
生产事故面试重点。
20. 方法1:增加消费者数量(优先)
例如:
Kubernetes:
原:
replicas: 3
调整:
replicas: 20
提高消费速度。
21. 方法2:临时跳过历史消息
如果消息已经无意义:
例如:
实时日志。
修改 offset:
kafka-consumer-groups.sh \
--reset-offsets \
--to-latest
例如:
跳过历史:
10亿条旧日志
⚠️ 会造成数据丢失,需要业务确认。
22. 方法3:增加 Kafka retention 时间
如果积压严重:
避免消息被删除。
查看:
kafka-configs.sh \
--describe
调整:
retention.ms
23. 方法4:临时扩容消费者资源
Kubernetes:
增加:
CPU
Memory
例如:
resources:
requests:
cpu: 4
24. 生产环境完整处理流程(面试高级回答)
面试可以这样说:
Kafka 消息积压首先通过 consumer group lag 判断积压情况,然后确认积压发生在哪些 partition。接着分析是消费端处理能力不足、消费者异常、分区设计问题还是 Kafka broker 性能问题。如果是消费能力不足,会优先增加 consumer 实例或者提高消费并发;如果是业务处理慢,则优化消费逻辑,例如批量处理、异步化。如果是 partition 不足,会扩容 partition。如果是 broker 瓶颈,需要扩容 broker 或优化磁盘网络。紧急情况下,可以临时扩容消费者或者根据业务允许调整 offset 跳过无效消息。
25. K8s环境下 Kafka 积压排查
如果 Kafka 跑在 Kubernetes:
25.1. 查看消费者 Pod
kubectl get pod -n kafka
25.2. 查看资源
kubectl top pod
25.3. 查看日志
kubectl logs consumer-pod
25.4. 扩容消费者
kubectl scale deployment consumer \
--replicas=20
26. 面试追问:
面试官通常还会继续问:
-
Kafka 为什么增加消费者不能无限提升性能?
-
Kafka Consumer Group 如何保证消息不重复消费?
-
Kafka 如何保证消息不丢失?
-
Kafka Lag 很高但是消费者 CPU 不高,可能是什么原因?
-
Kafka Rebalance 为什么会导致积压?
-
如何监控 Kafka 消息积压?
这些是 K8s + Kafka 运维岗(2~5年经验)非常常见的深入问题。建议下一步重点准备 Kafka 消息丢失、重复消费、Exactly Once 语义,这几个和积压经常连着问。
