ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

用 Kafka 做订单超时关闭,我们重写了三遍:Kafka 与 RocketMQ 的 5 个决定性差异

用 Kafka 做订单超时关闭,我们重写了三遍:Kafka 与 RocketMQ 的 5 个决定性差异 sn: 20batch: 4round: 9topic: 消息队列 Kafka 与 RocketMQ 选型对比订单 30 分钟未支付自动关闭这个需求听起来人畜无害。我们第一个版本用 Kafka 做的下单发消息消费端 sleep不行。定时轮询订单表delay 不精确。最后妥协成Kafka 定时扫表补偿代码绕了三层。两年后另一个团队接手新业务直接上了 RocketMQ 的延迟消息一个 send 方法搞定代码量是我们的三分之一。那次技术评审上我承认当初选 Kafka 不是技术判断是团队熟。这篇把两个 MQ 在架构、延迟消息、事务消息、消息顺序、运维成本五个维度的差异讲透附上我压测过的真实数据和选型决策树。版本Kafka 3.7RocketMQ 5.2。先看架构差异这决定了它们天生的性格Kafka 的 topic 是逻辑概念物理上是分区partition分散在多台 broker 上没有主从复制角色KRaft 模式后 broker 自治。RocketMQ 是经典的主从架构NameServer 做轻量路由broker 分主从。架构差异直接投射到使用体验上Kafka 的顺序性绑定在分区级别一个订单的毫秒级延迟消息、事务消息这类消息功能Kafka 社区要么没有要么很别扭RocketMQ 把这些做成一等公民功能。反过来Kafka 的大吞吐日志流场景埋点、CDC、日志汇聚是 RocketMQ 短期追不上的。延迟消息这是最硬的差距RocketMQ 5.x 支持任意精度的延迟消息4.x 是 18 个固定等级。订单超时关闭的代码是这样的Message msg new Message( ORDER_TIMEOUT_TOPIC, // 1. 目标 topic close, // 2. tag orderId.getBytes(StandardCharsets.UTF_8)); msg.setDeliverTimeMs(System.currentTimeMillis() 30 * 60 * 1000); // 3. 30 分钟后投递 SendResult result producer.send(msg); if (result.getSendStatus() ! SendStatus.SEND_OK) { // 4. 发送失败走本地重试或落库补偿 timeoutCompensator.save(orderId, deliverAt); }逐行看第 3 行设置精确投递时间broker 会把消息按投递时间放入内部的 timer wheel5.x 的 TimerWheel 实现时间轮 多级滚动到点再转入目标 topic第 4 行很重要发送失败必须落库补偿MQ 的 send 成功才是订单超时逻辑的起点失败就意味着这个订单不会被自动关闭。Kafka 要实现同样效果主流方案三种定时轮询 DB、自己搭延迟队列服务内部用 Kafka Redis ZSet、或者用 Kafka Streams 的 punctuate。三种都是拿复杂度换功能。我们把旧系统从 Kafka 延迟方案迁到 RocketMQ 后删掉了 1400 行轮询代码超时关闭的误差从轮询间隔 1 分钟降到 100 毫秒内。事务消息RocketMQ 的杀手锏Kafka 要靠事务语义变通分布式事务的本地消息表方案用 RocketMQ 事务消息可以省掉本地表。发送方先发 half 消息执行本地事务然后 commit 或 rollbackproducer.sendMessageInTransaction(msg, new TransactionListener() { Override public LocalTransactionState executeLocalTransaction(Message msg, Object arg) { try { // 1. 执行本地事务扣减库存 stockService.deduct(parseOrderId(msg)); // 2. 本地事务成功提交 half 消息消费者才能看到 return LocalTransactionState.COMMIT_MESSAGE; } catch (Exception e) { // 3. 本地事务失败回滚 half 消息消息作废 return LocalTransactionState.ROLLBACK_MESSAGE; } } Override public LocalTransactionState checkLocalTransaction(MessageExt msg) { // 4. broker 回查half 消息悬而未决时比如发送方宕机检查本地事务状态 boolean done stockService.isDeducted(parseOrderId(msg)); return done ? LocalTransactionState.COMMIT_MESSAGE : LocalTransactionState.ROLLBACK_MESSAGE; } });逐行拆第 1-2 行本地事务和消息提交绑在一个回调里half 消息在 commit 前对消费者不可见broker 存在系统内部 topic第 4 行的回查机制是这个设计的精华——发送方执行本地事务后宕机half 消息悬着broker 会按间隔反复回查默认 60 秒一次最多 15 次直到得到明确答复。注意回查逻辑必须幂等且能从 DB 还原事实不能依赖内存状态。Kafka 的事务是另一种东西它解决的是多条消息跨分区原子写入 消费端 read_committed本质是流处理场景的 exactly-once 支持不是业务事务消息。拿 Kafka 事务模拟本地事务 消息语义得自己实现回查逻辑工作量不小还容易错。压测数据与选型决策树同一台 16C64G 的三节点集群我们的实测Kafka 单 topic 100 分区生产 50 万 msg/s1KB 消息无堆积CPU 还有余量RocketMQ 同硬件规模生产 25-30 万 msg/s。但 RocketMQ 开启延迟消息和事务消息后功能复杂度换来的吞吐下降可控。吞吐量级差异的真实原因在存储设计Kafka 顺序追加 零拷贝 sendfile 把网络发送做到了极致RocketMQ 的 CommitLog 统一存储 ConsumeQueue 索引设计在功能扩展上更灵活。我的决策树供参考你的场景我的推荐理由日志汇聚、埋点、CDC 管道Kafka吞吐天花板、生态Flink/Connect成熟订单/交易/延迟/事务消息RocketMQ延迟与事务是一等公民团队只有 3-5 人维护中间件二选一即可两个都稳别同时引入需要消费重试 死信队列开箱即用RocketMQKafka 的 DLQ 要自己拼最后说个观点我见过团队为了统一技术栈把延迟消息硬塞进 Kafka也见过为了功能全用 RocketMQ 扛日均 TB 级日志流结果都是运维地狱。选型的第一问不是哪个更强而是你的消息里有多少业务语义。语义多延迟、事务、顺序、重试RocketMQ语义少、吞吐大Kafka。消费端的真实差距重试队列与死信是内置的选型时大家盯生产端吞吐实际运维里吃掉最多工时的是消费端异常处理。RocketMQ 的重试 死信是开箱即用的consumer.subscribe(ORDER_PAID_TOPIC, *); consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) - { for (MessageExt msg : msgs) { try { handleOrderPaid(msg); // 1. 处理成功消费确认 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } catch (BizException e) { // 2. 业务异常返回 RECONSUME_LATERbroker 把消息挪进重试队列 // 重试间隔按 %RETRY%topic 消费组名 的延迟级别递增10s/30s/1m/2m... log.warn(biz fail, reconsume later, msgId{}, msg.getMsgId(), e); return ConsumeConcurrentlyStatus.RECONSUME_LATER; } } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; });逐行说第 2 行的重试不是简单的重新投递——broker 维护%RETRY%前缀的重试 topic每次重试间隔按延迟级别递增最多 16 次后进%DLQ%死信 topic等人工介入。这个指数退避 死信兜底的设计Kafka 全部要自己做重试 topic 建几个、间隔怎么排、死信存哪里、消费位点怎么回拨每一条都是自研代码和运维流程。我们把 Kafka 消费者的重试逻辑迁到 RocketMQ 后删了大约 600 行重试编排代码。再补一个顺序消费的对比。订单状态机要求同一订单的消息按序消费RocketMQ 用 MessageQueueSelector 顺序监听器// 发送端按 orderId 哈希选队列同订单消息固定进同一队列 producer.send(msg, (mqs, message, arg) - { long orderId (long) arg; int index (int) (Math.abs(orderId) % mqs.size()); return mqs.get(index); }, orderId); // 消费端MessageListenerOrderly 保证队列内串行 队列级锁 consumer.registerMessageListener((MessageListenerOrderly) (msgs, context) - { // 同一队列的消息串行执行且队列被当前线程独占 for (MessageExt msg : msgs) { applyStatusTransition(msg); // 状态机按序演进 } return ConsumeOrderlyStatus.SUCCESS; });逐行拆发送端的哈希选队列保证同 key 消息落同一队列分区有序的前提消费端的 Orderly 监听器对队列加分布式锁队列内严格串行——代价是没有重试退避失败会原地阻塞重试所以顺序消费的业务必须把异常分类可重试的继续抛、不可重试的吞掉记日志否则一条毒消息会堵住整个队列。Kafka 要达到同样语义代码形态差不多按 key 分区 单分区单线程但它没有顺序失败原地阻塞的现成机制乱序保护弱一些。思考题RocketMQ 5.x 的定时消息基于时间轮如果允许设置一年后的投递时间时间轮的内存和磁盘会怎样为什么 RocketMQ 官方把默认最大延迟限制在 3 天评论区聊聊你的推断。
返回列表