Kafka 消费组 rebalance 风暴堆积 200 万条那晚:和 RocketMQ 比,5 个真正决定选型的差异
2026/8/5 13:10:54 网站建设 项目流程

title: Kafka 消费组 rebalance 风暴堆积 200 万条那晚:和 RocketMQ 比,5 个真正决定选型的差异
tags: [Kafka, RocketMQ, 消息队列, 中间件选型, Java]
category: 后端


一次发布,把消费组拖进了死循环

那晚是常规发版,滚动重启 12 个消费实例。按经验这事 3 分钟就能结束,结果监控上的 lag 曲线一路往上冲,20 分钟涨到 200 万条,消费速率几乎归零。

登上机器看日志,满屏都是这个:

[Consumer clientId=order-consumer-7, groupId=order-group] Attempt to heartbeat failed since group is rebalancing [Consumer clientId=order-consumer-7, groupId=order-group] Revoke previously assigned partitions order-topic-3, order-topic-11 [Consumer clientId=order-consumer-7, groupId=order-group] (Re-)joining group

rebalance 一轮接一轮,永远结束不了。这就是所谓的 rebalance 风暴。

根因是三个配置叠在一起:

  1. max.poll.interval.ms用的默认值 300000(5 分钟),但我们单条消息的处理逻辑里有个外部 HTTP 调用没设超时,偶尔会卡 6 分钟以上
  2. max.poll.records是默认 500,一次拉 500 条,只要有几条卡住,整批就处理不完
  3. 滚动重启用的是默认的RangeAssignor,每次有实例进出,全部分区都要重新分配

于是形成了闭环:实例 A 处理超时被踢出组 → 触发 rebalance → 所有实例暂停消费重新分配 → 分配完 A 又拉了一批带毒消息 → 再次超时被踢 → 再 rebalance。

那晚的处理办法很粗暴:把消费组的group.instance.id加上(启用静态成员),max.poll.records从 500 降到 50,给那个 HTTP 调用加了 2 秒超时,然后重启。lag 在 40 分钟后清零。

复盘会上有人问:「如果我们用的是 RocketMQ,还会有这个问题吗?」这个问题让我把两边的消费模型认真对比了一遍。这篇写的就是这次对比的结论。环境是 Kafka 3.2.1(12 分区)、RocketMQ 4.9.4、Spring Boot 2.7.5、JDK 11。

差异一:消费模型和 rebalance 的处理方式

Kafka 的 rebalance 由 Broker 端的 GroupCoordinator 主导,走的是 JoinGroup / SyncGroup 两阶段协议。关键特征是Stop-The-World:rebalance 期间整个消费组停止消费。

RocketMQ 的 rebalance 是客户端各自算的。每个 Consumer 定时(默认 20 秒)从 NameServer 拉取 Topic 路由和消费组成员列表,然后用相同的算法(默认AllocateMessageQueueAveragely)独立计算自己该消费哪些队列。

// RocketMQ RebalanceImpl#rebalanceByTopic 的核心(简化) List<MessageQueue> mqAll = new ArrayList<>(mqSet); Collections.sort(mqAll); // 队列排序,保证所有客户端看到一致的顺序 Collections.sort(cidAll); // 消费者 ID 排序,同上 AllocateMessageQueueStrategy strategy = this.allocateMessageQueueStrategy; // 每个客户端独立计算,因为输入相同、算法相同,结果必然一致 List<MessageQueue> allocateResult = strategy.allocate( this.consumerGroup, this.mQClientFactory.getClientId(), mqAll, cidAll); Set<MessageQueue> allocateResultSet = new HashSet<>(allocateResult); // 只更新有变化的队列,没变化的队列消费完全不受影响 boolean changed = this.updateProcessQueueTableInRebalance(topic, allocateResultSet, isOrder);

updateProcessQueueTableInRebalance这个方法是关键:它做的是增量更新——对比新旧分配结果,只丢弃不再属于自己的队列,只新增分配给自己的队列。没有变化的队列,消费线程压根不知道发生过 rebalance。

这个差异在实践中的影响:

场景KafkaRocketMQ
滚动重启 12 个实例每次实例进出触发全组 STW,12 次 rebalance只有涉及的队列被迁移,其余照常消费
单个消费者卡死max.poll.interval超时后被踢,全组 rebalance该消费者的队列会被重新分配,其他不受影响
扩容一个实例全组 STW 一次部分队列迁移

Kafka 从 2.3 开始有了CooperativeStickyAssignor(增量协作式再平衡),能大幅减少 STW 范围。我们那次事故后就换成了它,配合group.instance.id静态成员,滚动重启期间的消费中断从「全组停 20 分钟」降到「单实例停 3 秒」。

props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, CooperativeStickyAssignor.class.getName()); // 静态成员:实例重启后用同一个 group.instance.id 重新加入, // 在 session.timeout.ms 内不触发 rebalance props.put(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, "order-consumer-" + podOrdinal); props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 45000); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 50); props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 120000);

group.instance.id用 Pod 序号而不是随机 UUID,这点很重要。K8s 里如果用 Deployment,Pod 名是随机的;我们改成了 StatefulSet,Pod 名固定为order-consumer-0order-consumer-11,重启后 ID 不变,才能享受静态成员的好处。

差异二:延迟消息,一个内置一个要自己造

RocketMQ 原生支持 18 个延迟级别:

Message msg = new Message("order-timeout-topic", body); // level 16 = 30 分钟,用于订单超时未支付自动关单 msg.setDelayTimeLevel(16); producer.send(msg);

实现原理是 Broker 把带延迟级别的消息先写进内部 TopicSCHEDULE_TOPIC_XXXX,每个延迟级别对应一个队列,用定时任务扫描到期消息再投递到真实 Topic。

// ScheduleMessageService$DeliverDelayedMessageTimerTask 的核心逻辑 long now = System.currentTimeMillis(); long deliverTimestamp = computeDeliverTimestamp(delayLevel, storeTimestamp); long countdown = deliverTimestamp - now; if (countdown > 0) { // 还没到时间,重新调度自己,间隔 100ms this.scheduleNextTimerTask(nextOffset, DELAY_FOR_A_WHILE); return; } // 到期了:把消息从 SCHEDULE_TOPIC 恢复成原 Topic 再投递 MessageExtBrokerInner msgInner = messageTimeup(msgExt); PutMessageResult result = defaultMessageStore.putMessage(msgInner);

固定 18 个级别(1s / 5s / 10s / 30s / 1m / 2m / 3m / 4m / 5m / 6m / 7m / 8m / 9m / 10m / 20m / 30m / 1h / 2h)是它的限制。想要「延迟 47 分钟」就得自己组合或者改 Broker 配置。RocketMQ 5.0 引入了基于时间轮的任意精度定时消息,但我们生产上还是 4.9.4,没用上。

Kafka 完全没有延迟消息。要实现就得自己搭:常见做法是建 N 个延迟 Topic(delay-5sdelay-30sdelay-5m...),消费者拉到消息后判断是否到期,没到期就pause()分区并seek()回去。

// Kafka 实现延迟消费的典型写法 ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500)); for (ConsumerRecord<String, String> record : records) { long deliverAt = Long.parseLong(new String(record.headers().lastHeader("deliverAt").value())); if (System.currentTimeMillis() < deliverAt) { // 没到时间:暂停这个分区,并把 offset 拨回这条消息 TopicPartition tp = new TopicPartition(record.topic(), record.partition()); consumer.pause(Collections.singleton(tp)); consumer.seek(tp, record.offset()); pausedUntil.put(tp, deliverAt); break; // 同一分区后面的消息也不用看了,因为投递时间是递增的 } process(record); } // 定期检查是否该 resume pausedUntil.entrySet().removeIf(e -> { if (System.currentTimeMillis() >= e.getValue()) { consumer.resume(Collections.singleton(e.getKey())); return true; } return false; });

这段代码能跑,但有几个问题:break依赖「同一分区内投递时间递增」这个假设,生产者乱序发送就失效了;pause期间该分区完全不消费,如果分区里混了不同延迟时长的消息就会阻塞;每个延迟档位都要独立 Topic 和消费者,运维复杂度上去了。

如果你的业务大量依赖延迟消息(订单超时、定时提醒、重试退避),我不建议选 Kafka。自己造这套轮子的成本和后续维护成本都不低。

差异三:消息重试和死信,一个自动一个手动

RocketMQ 消费失败后返回RECONSUME_LATER,Broker 会把消息投递到重试 Topic%RETRY%{consumerGroup},按延迟级别递增重试 16 次,全部失败后进入死信 Topic%DLQ%{consumerGroup}。整套流程零代码。

@Override public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) { MessageExt msg = msgs.get(0); try { orderService.handle(msg); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } catch (RetryableException e) { // 返回这个,Broker 自动安排重试,第 N 次重试的延迟级别是 N+2 return ConsumeConcurrentlyStatus.RECONSUME_LATER; } catch (Exception e) { // 不可重试的异常,直接消费成功(丢弃)+ 记录告警,避免无意义重试 alarmService.warn("unrecoverable msg, msgId=" + msg.getMsgId(), e); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }

getReconsumeTimes()能直接拿到当前是第几次重试,做「重试三次后转人工」这类逻辑很方便。

Kafka 这边什么都没有。Spring Kafka 的SeekToCurrentErrorHandler(新版本叫DefaultErrorHandler)提供了本地重试,但它是阻塞式的——重试期间这个分区的后续消息全部堵住。生产上更常见的做法是自建重试 Topic 链:

@Bean public DefaultErrorHandler errorHandler(KafkaTemplate<String, String> template) { // 失败的消息发到 <原topic>-retry,重试 3 次后进 <原topic>-dlt DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(template, (record, ex) -> { int attempts = getAttempts(record); String target = attempts >= 3 ? record.topic() + "-dlt" : record.topic() + "-retry"; return new TopicPartition(target, -1); }); // 本地不重试,直接转发,避免阻塞分区 return new DefaultErrorHandler(recoverer, new FixedBackOff(0L, 0L)); }

new FixedBackOff(0L, 0L)表示本地零重试。这是我们踩过坑之后的选择——最初用的是FixedBackOff(1000L, 3L),一条毒消息会让整个分区停 3 秒,高峰期直接把 lag 堆起来。

差异四:顺序消息的粒度

Kafka 保证的是分区内有序。要让同一订单的消息有序,就得让它们落到同一分区,靠 key 的哈希实现。

RocketMQ 保证的是队列内有序,用MessageQueueSelector指定队列,并且消费端要用MessageListenerOrderly。区别在于 RocketMQ 的顺序消费会对队列加锁,同一队列同一时刻只有一个线程消费;Kafka 是一个分区对应一个消费者线程,天然串行。

实际差异在失败处理上。RocketMQ 顺序消费失败返回SUSPEND_CURRENT_QUEUE_A_MOMENT,会阻塞当前队列直到成功(默认最多Integer.MAX_VALUE次),保证严格顺序。Kafka 没有这个概念——你自己 commit offset 就跳过了,不 commit 就重复消费,没有中间态。

这一条我认为 RocketMQ 明确更强。但也要提醒:顺序消费失败会阻塞整个队列,如果毒消息永远处理不成功,这个队列就永久卡住了。我们的做法是加一个「顺序消费失败超过 100 次转异步兜底队列」的逻辑,牺牲严格顺序换可用性。

五个维度的综合对比

维度Kafka 3.2RocketMQ 4.9我们的判断
吞吐(单机,1KB 消息)我们压测约 78 万 TPS约 32 万 TPSKafka 明显更强
Rebalance 影响STW(3.x 有 Cooperative 缓解)增量,影响局部RocketMQ 更平滑
延迟消息无,需自建18 级内置,5.0 支持任意精度RocketMQ 完胜
消息重试/死信需自建 Topic 链内置 16 次重试 + DLQRocketMQ 完胜
事务消息有,但只保证「生产者到 Broker」半消息 + 回查,覆盖本地事务RocketMQ 更完整
消息回溯按 offset 或时间戳 seek按时间戳回溯打平
消息过滤消费端过滤Broker 端 Tag/SQL92 过滤RocketMQ 省带宽
生态Connect / Streams / ksqlDB 极丰富相对单薄Kafka 完胜
运维复杂度ZK 依赖(KRaft 后减轻)、分区规划难NameServer 无状态,简单RocketMQ 更省心

我们最后是怎么分的

没有全部切换,而是按场景拆开:

留在 Kafka 的:埋点日志(日均 40 亿条)、用户行为流、数据同步 CDC。这些场景吞吐要求极高、对延迟消息和重试机制没需求、下游还要接 Flink 做实时计算——Kafka 的生态优势在这里无可替代。

迁到 RocketMQ 的:订单超时关单(需要延迟消息)、支付回调(需要重试 + 死信)、库存扣减(需要事务消息)。这些是业务链路,消息量不大(日均 3000 万级),但对可靠性和功能完整性要求高。

迁移过程中最花时间的不是代码,是幂等。Kafka 和 RocketMQ 都只保证 at-least-once,重复投递是常态。我们统一用「消息 ID + Redis SETNX + 业务表唯一索引」三层去重,这套逻辑抽成了一个 starter,两边共用。

如果一定要给个一句话建议:日志流、数据管道选 Kafka;业务消息、需要延迟和事务的场景选 RocketMQ。至于「用一套统一技术栈」的诉求,我理解但不太认同——为了统一而在错误的场景上硬凑,后期补的轮子比省下来的运维成本贵得多。

顺便说一句,我们也评估过 Pulsar。存算分离的架构确实优雅,多租户和跨地域复制是亮点,但团队没人有生产运维经验,BookKeeper 那一层出问题不好排查,最后放弃了。选型除了技术指标,团队的掌控能力是一个不能忽略的权重项。

三个可以想想的问题

  1. Kafka 换成CooperativeStickyAssignor之后,从旧的RangeAssignor滚动升级需要两轮重启(先加入新策略,再移除旧策略)。为什么不能一次性切换?
  2. RocketMQ 的顺序消息在 Broker 主从切换时还能保证顺序吗?如果不能,业务侧要怎么兜底?
  3. 如果你的业务只需要「延迟 15 分钟」这一个档位,用 Kafka 自建延迟 Topic 和引入 RocketMQ,成本上你会怎么算这笔账?

如果你正在被 rebalance 折磨,先去看两个配置:partition.assignment.strategygroup.instance.id。这两项改完,八成的问题会自己消失。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询