Kafka消息堆积应该如何处理?
消息堆积处理建议当线上发生严重积压影响下游业务时需要根据分区数与消费者数量的关系采取对应策略方案 1消费者数 分区数直接扩容 ConsumerKafka 中一个分区同一时间只能被同一个 Group 内的一个 Consumer 消费。若消费节点数未达分区上限直接横向扩容。方案 2消费者数 分区数临时重路由当消费者数量已达分区数上限时直接加 Consumer 是无效的。此时需通过“中转队列”突破分区限制新建临时 Topic创建一个拥有更多分区如原分区的 10 倍的新 Topic。上线临时转存程序写一个不包含复杂业务逻辑的 Consumer专门从原 Topic 批量拉取积压消息写入到新的 Topic 中部署多倍 Consumer部署 10 倍于原规模的 Consumer 实例挂载到新 Topic 上并行消费清理恢复积压清空后切回原有的消费链路并下线临时设施方案 3紧急丢弃 / 转存 DB非核心业务丢弃非关键日志若堆积的是日志或时效性要求极高的监控数据且业务允许丢失可直接通过重置消费位点Offset到最新位置Latest快速跳过积压。写入数据库若消息不能丢但消费太慢可以先将积压消息批量 Dump 到 MySQL / Redis / ES 中再由离线批处理任务慢慢消化。方案4跳过坏消息如果某条消息的消费逻辑有问题导致消费异常并且消息消费无顺序性要求时可以先跳过这条消息排查根因与长期防范堆积通常由消费端卡死、下游性能瓶颈或生产者突发大流量导致。排查与优化策略如下1. 排查消费端阻塞与异常查看线程 Dump检查 Consumer 进程是否因调用第三方 API 超时、数据库死锁、慢 SQL 等导致消费线程卡死。优化单条消费逻辑将 RPC 串行调用改为异步/并行调用将单条写 DB 改为批量写入如 Batch Insert。2. 调整 Kafka 消费参数参数说明建议配置 / 优化方向max.poll.records每次 poll 拉取的最大消息数若单条处理极慢调小该值防止触发消费超时若单条处理快调大以提升吞吐。max.poll.interval.ms两次 poll 的最大间隔时间若业务处理耗时较长增加此时间避免被误判为宕机而频繁触发 Rebalance。fetch.min.bytes/fetch.max.wait.ms批次拉取阈值调高可让 Consumer 一次拉取更多数据提升网络传输吞吐率。3. 优化生产端与 Topic 架构合理设置分区数创建 Topic 时预留足够的 Partition 数如根据峰值 TPS 计算为未来的弹性扩容留出空间。 如果后期评估分区数量不足增加 topic 分区数同步扩容消费topic的分区数可以动态增加但不能动态减少打散消息避免生产者固定 key 导致消息倾斜集中在个别分区非顺序性业务可不指定 key 或追加随机后缀