尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

Kafka 消费组 rebalance 风暴堆积 200 万条那晚:和 RocketMQ 比,5 个真正决定选型的差异

Kafka 消费组 rebalance 风暴堆积 200 万条那晚:和 RocketMQ 比,5 个真正决定选型的差异 title: Kafka 消费组 rebalance 风暴堆积 200 万条那晚和 RocketMQ 比5 个真正决定选型的差异tags: [Kafka, RocketMQ, 消息队列, 中间件选型, Java]category: 后端一次发布把消费组拖进了死循环那晚是常规发版滚动重启 12 个消费实例。按经验这事 3 分钟就能结束结果监控上的 lag 曲线一路往上冲20 分钟涨到 200 万条消费速率几乎归零。登上机器看日志满屏都是这个[Consumer clientIdorder-consumer-7, groupIdorder-group] Attempt to heartbeat failed since group is rebalancing [Consumer clientIdorder-consumer-7, groupIdorder-group] Revoke previously assigned partitions order-topic-3, order-topic-11 [Consumer clientIdorder-consumer-7, groupIdorder-group] (Re-)joining grouprebalance 一轮接一轮永远结束不了。这就是所谓的 rebalance 风暴。根因是三个配置叠在一起max.poll.interval.ms用的默认值 3000005 分钟但我们单条消息的处理逻辑里有个外部 HTTP 调用没设超时偶尔会卡 6 分钟以上max.poll.records是默认 500一次拉 500 条只要有几条卡住整批就处理不完滚动重启用的是默认的RangeAssignor每次有实例进出全部分区都要重新分配于是形成了闭环实例 A 处理超时被踢出组 → 触发 rebalance → 所有实例暂停消费重新分配 → 分配完 A 又拉了一批带毒消息 → 再次超时被踢 → 再 rebalance。那晚的处理办法很粗暴把消费组的group.instance.id加上启用静态成员max.poll.records从 500 降到 50给那个 HTTP 调用加了 2 秒超时然后重启。lag 在 40 分钟后清零。复盘会上有人问「如果我们用的是 RocketMQ还会有这个问题吗」这个问题让我把两边的消费模型认真对比了一遍。这篇写的就是这次对比的结论。环境是 Kafka 3.2.112 分区、RocketMQ 4.9.4、Spring Boot 2.7.5、JDK 11。差异一消费模型和 rebalance 的处理方式Kafka 的 rebalance 由 Broker 端的 GroupCoordinator 主导走的是 JoinGroup / SyncGroup 两阶段协议。关键特征是Stop-The-Worldrebalance 期间整个消费组停止消费。RocketMQ 的 rebalance 是客户端各自算的。每个 Consumer 定时默认 20 秒从 NameServer 拉取 Topic 路由和消费组成员列表然后用相同的算法默认AllocateMessageQueueAveragely独立计算自己该消费哪些队列。// RocketMQ RebalanceImpl#rebalanceByTopic 的核心简化 ListMessageQueue mqAll new ArrayList(mqSet); Collections.sort(mqAll); // 队列排序保证所有客户端看到一致的顺序 Collections.sort(cidAll); // 消费者 ID 排序同上 AllocateMessageQueueStrategy strategy this.allocateMessageQueueStrategy; // 每个客户端独立计算因为输入相同、算法相同结果必然一致 ListMessageQueue allocateResult strategy.allocate( this.consumerGroup, this.mQClientFactory.getClientId(), mqAll, cidAll); SetMessageQueue allocateResultSet new HashSet(allocateResult); // 只更新有变化的队列没变化的队列消费完全不受影响 boolean changed this.updateProcessQueueTableInRebalance(topic, allocateResultSet, isOrder);updateProcessQueueTableInRebalance这个方法是关键它做的是增量更新——对比新旧分配结果只丢弃不再属于自己的队列只新增分配给自己的队列。没有变化的队列消费线程压根不知道发生过 rebalance。这个差异在实践中的影响场景KafkaRocketMQ滚动重启 12 个实例每次实例进出触发全组 STW12 次 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 里如果用 DeploymentPod 名是随机的我们改成了 StatefulSetPod 名固定为order-consumer-0到order-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 个延迟 Topicdelay-5s、delay-30s、delay-5m...消费者拉到消息后判断是否到期没到期就pause()分区并seek()回去。// Kafka 实现延迟消费的典型写法 ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(500)); for (ConsumerRecordString, 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_LATERBroker 会把消息投递到重试 Topic%RETRY%{consumerGroup}按延迟级别递增重试 16 次全部失败后进入死信 Topic%DLQ%{consumerGroup}。整套流程零代码。Override public ConsumeConcurrentlyStatus consumeMessage(ListMessageExt msgs, ConsumeConcurrentlyContext context) { MessageExt msg msgs.get(0); try { orderService.handle(msg); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } catch (RetryableException e) { // 返回这个Broker 自动安排重试第 N 次重试的延迟级别是 N2 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(KafkaTemplateString, 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 影响STW3.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 那一层出问题不好排查最后放弃了。选型除了技术指标团队的掌控能力是一个不能忽略的权重项。三个可以想想的问题Kafka 换成CooperativeStickyAssignor之后从旧的RangeAssignor滚动升级需要两轮重启先加入新策略再移除旧策略。为什么不能一次性切换RocketMQ 的顺序消息在 Broker 主从切换时还能保证顺序吗如果不能业务侧要怎么兜底如果你的业务只需要「延迟 15 分钟」这一个档位用 Kafka 自建延迟 Topic 和引入 RocketMQ成本上你会怎么算这笔账如果你正在被 rebalance 折磨先去看两个配置partition.assignment.strategy和group.instance.id。这两项改完八成的问题会自己消失。
返回列表