Kafka

Kafka 消息积压处理

·9 分钟阅读·3495 字

Kafka 消息积压的确认方法、原因分析与消费者侧优化方案

📋 目录

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. 面试追问:

面试官通常还会继续问:

  1. Kafka 为什么增加消费者不能无限提升性能?

  2. Kafka Consumer Group 如何保证消息不重复消费?

  3. Kafka 如何保证消息不丢失?

  4. Kafka Lag 很高但是消费者 CPU 不高,可能是什么原因?

  5. Kafka Rebalance 为什么会导致积压?

  6. 如何监控 Kafka 消息积压?

这些是 K8s + Kafka 运维岗(2~5年经验)非常常见的深入问题。建议下一步重点准备 Kafka 消息丢失、重复消费、Exactly Once 语义,这几个和积压经常连着问。

Yanche Blog

记录云原生、Linux、数据库等技术领域的学习心得,以及日常生活的思考与感悟。

© 2026 Yanche Blog. All rights reserved.

Powered by Astro