Kafka积压怎么排查?消息队列选型别只看吞吐:RabbitMQ、Pulsar、RocketMQ、Kafka 踩坑对比

🔑 关键词:消息队列选型,Kafka积压排查,RabbitMQ,Pulsar,RocketMQ

📖 摘要:从一次80万日订单、峰值1200 TPS的Kafka积压事故讲起,对比Kafka/RabbitMQ/Pulsar/RocketMQ在顺序、事务、积压、扩容上的真实差异,给出可落地的排查步骤和参数。

先交代背景:80万日订单,Kafka 积压 400 万,我差点把消费者扩到 100 个

图片

去年双十一前,我们订单系统日均 80 万单,峰值 1200 TPS,平均 300 TPS。Kafka 是 3 台 8C16G,12 个分区,副本 2,acks=1,batch.size=16384,linger.ms=5。大促当天凌晨,lag 从 2000 冲到 400 万,消费者 8 个,max.poll.records=500,但每个消费者内部是逐条写 MySQL,单条 20ms,一批 500 条要 10 秒,算下来 8 个消费者只有 400 条/s。我第一反应是扩消费者,从 8 扩到 24,lag 没降,因为 MySQL 连接池只有 20,线程池 8,DB CPU 已经 85%,每个事务都在等锁。最后把写库改成 500 条一次 insert,rewriteBatchedStatements=true,连接池 20 到 50,消费者降到 12,lag 从 400 万降到 3 万用了 11 分钟。这个事让我明白:消息队列积压不是消费者慢,是生产速度减消费速度的积分,先看下游,再扩消费者。

我的独立观点:选型先问三个问题,别问吞吐

图片

很多文章一上来就贴 Kafka 百万吞吐、RabbitMQ 低延迟,我觉得容易把人带沟里。你先问自己三个问题:这条消息能不能丢?能不能重复?能不能乱序?订单支付消息不能丢、不能乱,重复要靠幂等;日志采集可以丢一点、可以乱,但吞吐要大。答案不同,选型完全不同。我现在的原则是:消息队列不是数据库,也不是业务状态机本身,它只是把状态变更搬来搬去的日志。积压是成本,磁盘、内存、恢复时间都是钱;lag 不是唯一指标,lag 除以消费速率才是恢复时间。小团队先治理生产者:限流、幂等、重试、死信,消费者扩容是最贵、最晚做的事。

四个消息队列的横向对比:我压测和踩坑后的表

图片

下面数据不是官方 benchmark,是我自己用 3 台 8C16G 云主机、1KB 消息、3 副本/3 节点压出来的大概值,别当标准,当参考。Kafka 3 分区 1 副本,producer 约 7.8 万条/s,p99 12ms;RabbitMQ quorum queue 3 节点,约 1.2 万条/s,p99 45ms;Pulsar 3 bookie,约 5.6 万条/s,p99 18ms;RocketMQ 4.9 3 节点,约 6.1 万条/s,p99 15ms。

维度 Kafka RabbitMQ Pulsar RocketMQ
核心模型 分区日志 队列/交换机 Topic+分区+BookKeeper CommitLog+ConsumeQueue
顺序 分区内有序 队列内有序 Key_Shared 可按键有序 队列内有序
事务 0.11+ 幂等和事务 publisher confirm,不做分布式事务 支持事务 事务消息+回查
积压能力 磁盘顺序写,强 内存压力大,lazy/quorum 缓解 存算分离,强 磁盘,强
扩容 分区可增不可减,消费者不能超过分区 队列可增,镜像队列麻烦 broker 和 bookie 分开扩 队列可扩,NameServer 无状态
延迟 毫秒到几十毫秒 微秒到毫秒,低 毫秒 毫秒
运维 中,ZK/KRaft 要懂 低,但集群策略坑多 高,组件多 中,阿里系文档多
适合 日志、流、大数据 业务解耦、低延迟 多租户、混合负载 电商事务、顺序

图片

排查积压的 7 个步骤:我当时是这么干的

第一步,先看 lag 曲线和消费速率。命令是 kafka-consumer-groups --bootstrap-server 127.0.0.1:9092 --describe --group order-group,看 CURRENT-OFFSET、LOG-END-OFFSET、LAG。如果 lag 线性涨,生产大于消费;如果 lag 平,消费者卡死。第二步,jstack 看消费者线程栈,有没有卡在 MySQL socketRead、Redis、HTTP。第三步,看下游连接池 active/idle、DB 慢 SQL、Redis 热 key,我们当时就是 DB 连接池只有 20。第四步,限流生产者或降级,非核心消息扔到重试 topic,核心订单走同步。第五步,扩消费者前先算分区数,Kafka 12 分区扩到 24 消费者只有 12 个干活。第六步,批量写,MySQL 开 rewriteBatchedStatements=true,batch 500,连接池 50。第七步,监控里加恢复时间估算:lag 除以消费速率,超过 30 分钟就告警。

图片

关键参数别抄错:Kafka、RabbitMQ、Pulsar、RocketMQ 我常用的配置

Kafka 我一般设 acks=all,min.insync.replicas=2,unclean.leader.election.enable=false,max.poll.records=500,session.timeout.ms=10000,max.poll.interval.ms=300000。RabbitMQ 做延迟队列用 x-message-ttl=60000,死信用 x-dead-letter-exchange=dlx,优先级 x-max-priority=10;3.8 以后尽量用 Quorum Queue,3 副本,别再用镜像队列。Pulsar 我设 writeQuorum=3,ackQuorum=2,订阅用 Key_Shared,保留 7 天,bookie 至少 3 块 SSD。RocketMQ 事务消息会写 half topic,回查默认 15 次,broker 刷盘用 ASYNC_FLUSH,主从用 SYNC_MASTER,死信队列是 %DLQ%group。

图片

最后的选择建议,可能得罪人

日志、流、大数据,选 Kafka,别纠结。业务解耦、低延迟、中小规模,选 RabbitMQ,但一定用 Quorum Queue。多租户、存算分离、积压特别大,选 Pulsar,前提是你有人懂 BookKeeper。电商事务、顺序消息、阿里生态,选 RocketMQ。边缘和 IoT,可以看 NATS JetStream,R3 file storage,max_age 24h 就够。别把 MQ 当选型宗教,先问能不能丢、能不能重、能不能乱,再谈吞吐。顺序消息的代价是并行度,要顺序就按业务键 hash 到同一分区或队列,别全局顺序。事务消息也不是银弹,RocketMQ 回查 15 次也可能失败,最终一致还得靠本地事务表加定时补偿。

🏷️ 标签: