先说我遇到的事故:不是分区不够,是消费太慢
2023 年 10 月底,我们订单事件 topic order-event 积压到 180 万条。集群 3 个 broker,topic 12 分区,replication.factor=2,min.insync.replicas=1,acks=all。消费者组 order-consumer 有 6 个 Pod,每个 2 核 4G,JVM heap 2G,max.poll.records=500,enable.auto.commit=true。监控里 lag 从 20 万涨到 180 万,持续 40 分钟。我一开始以为加分区就行,后来发现日志每 5 分钟出现 Rebalance started,消费者频繁进出组。查代码,单批 500 条里每条调一次风控 HTTP,平均 320ms,最慢 1.8s。最坏一批 500*1.8s=900s,而默认 max.poll.interval.ms=300000,5 分钟。消费者被判定死亡,触发重平衡。这个坑很典型:分区数只决定并行上限,不决定单条处理速度。
我最后怎么调:参数和步骤
第一步,先别急着扩 broker。
用 kafka-consumer-groups.sh --bootstrap-server kafka1:9092 --describe --group order-consumer 看每个分区的 lag。如果 90% lag 集中在 2 个分区,那是 key 倾斜;如果每个分区都均匀涨,那是消费能力不够。我们的 12 个分区 lag 都在 8 万到 18 万之间,所以是后者。
第二步,把 max.poll.records 从 500 降到 100,max.poll.interval.ms 从 300000 调到 600000,session.timeout.ms=45000,heartbeat.interval.ms=3000。enable.auto.commit=false,改成每批处理完 commitSync()。
第三步,风控 HTTP 改成批量接口,100 条一次,平均 180ms,最慢 900ms。同时本地 Caffeine 缓存 5 分钟,减少 60% 调用。
第四步,数据库写入改成 rewriteBatchedStatements=true,每批 100 条,flush 间隔 500ms。
第五步,扩 broker 从 3 到 6,topic 分区从 12 到 24,replication.factor 改成 3,min.insync.replicas=2,acks=all。消费者 Pod 从 6 扩到 12,每个 Pod 2 个消费线程。注意:分区扩容后,新分区没有位点,auto.offset.reset 如果是 earliest,会把新分区从最早开始读,可能重复。我们业务有 Redis 去重,所以设 earliest;如果你们不能重复,先停消费者,用 --reset-offsets --to-latest 对新分区单独处理,别一把梭。
调整后 lag 从 180 万降到 2 万,用了 35 分钟。P99 消费延迟从 6 分钟降到 18 秒。
Kafka、Pulsar、RabbitMQ 我现在的选择逻辑
我以前会问“Kafka 能不能替代 RabbitMQ”。现在我不这么问。我先问三件事:失败消息要不要阻塞后面的消息?要不要按 key 保序?要不要回放历史?
Kafka 是分区提交日志。分区内有序,跨分区无序。消息被消费不等于删除,删除由 retention.ms 或 retention.bytes 控制。默认 retention.ms=604800000,7 天。它适合高吞吐、可回放、按 key 分区的场景。但它的位点控制权在消费者手里,自动提交就是“最多一次”或“至少一次”的混合体,别指望不重复。
Pulsar 是存算分离,broker 无状态,BookKeeper 存数据。subscription 是游标,支持 Exclusive、Failover、Shared、Key_Shared。Key_Shared 能让同一 key 落到同一消费者,但不同 key 可能共享消费者。它扩容快,但组件多:ZooKeeper 或 etcd、BookKeeper、broker。BookKeeper 一般 3 副本,ack quorum 2。运维复杂度比 Kafka 高,团队没专人别轻易上。
RabbitMQ 是队列模型,消费即删。Quorum Queue 基于 Raft,3 节点容忍 1 个挂,5 节点容忍 2 个挂。它适合低延迟、路由复杂、消息量在几千到几万 QPS 的场景。但队列多了,每个 Quorum Queue 都是一个 Raft 组,CPU 和内存开销明显。我们另一个项目 400 个 Quorum Queue,3 节点 8C16G,日常 CPU 35%,批量导入时到 70%。Kafka 同样量级只用了 6 个 broker 里的 2 个。
参数清单:我们生产现在用的
生产者:acks=all,enable.idempotence=true,min.insync.replicas=2,compression.type=lz4,linger.ms=10,batch.size=65536,max.in.flight.requests.per.connection=5,delivery.timeout.ms=120000,retries=2147483647。
消费者:enable.auto.commit=false,auto.offset.reset=earliest,max.poll.records=100,max.poll.interval.ms=600000,session.timeout.ms=45000,heartbeat.interval.ms=3000,fetch.min.bytes=1024,fetch.max.wait.ms=500,max.partition.fetch.bytes=1048576,isolation.level=read_committed。
Broker:num.partitions=24,default.replication.factor=3,min.insync.replicas=2,log.retention.hours=168,log.segment.bytes=1073741824,num.io.threads=8,num.network.threads=4,log.flush.interval.messages=10000,unclean.leader.election.enable=false。
版本:Kafka 3.6.1,KRaft 模式。Kafka 2.8 开始有 KRaft 早期访问,3.3 生产就绪,4.0 已经移除 ZooKeeper。如果你们还在 ZooKeeper,别慌,3.6 可以迁移,但先备份元数据。
我还建议用 kafka-producer-perf-test.sh 和 kafka-consumer-perf-test.sh 压测。我们 1KB 消息,acks=all,LZ4,3 副本,单分区写 8 MB/s 左右;24 分区理论 192 MB/s,实际受网卡和磁盘限制,跑到 110 MB/s。消费端单线程如果处理 10ms,单分区约 100 msg/s,24 分区 2400 msg/s。如果业务峰值 4500 msg/s,就得扩消费者线程或优化单条处理时间。
我的独立观点:别把 Kafka 当消息队列,把它当数据库日志
Kafka 最值钱的能力不是“发消息”,是“位点主权”和“回放”。你把位点交给自动提交,就等于把账本交给别人记。你把位点当成数据库事务的一部分,才能做 effectively-once。Kafka 事务只保证 Kafka 内部,跨 MySQL 要 outbox 表加幂等键,别迷信 exactly-once。 分区数也不是越多越好。分区多了,leader 选举、rebalance、文件句柄、端到端延迟都会涨。我们 6 个 broker 跑 240 个分区,每个 broker 40 个 leader,还能接受。如果到 6000 个分区,Controller 和元数据压力会明显上来。 选型时,如果业务是任务队列、延迟微秒级、路由到多个队列,RabbitMQ 更顺手。如果要多租户、存算分离、跨地域复制,Pulsar 值得看。如果是事件流、日志、回放、按 key 顺序、吞吐上万,Kafka 还是稳。但不管选哪个,先把“失败怎么办、重复怎么办、顺序怎么办”写在纸上,比对比 benchmark 有用。