尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

RocketMQ与Kafka:从存储模型、性能原理到业务选型

RocketMQ与Kafka:从存储模型、性能原理到业务选型 RocketMQ 与 Kafka从存储模型、性能原理到业务选型本文基于 RocketMQ 与 Kafka 的架构差异系统拆解 Kafka 为什么吞吐更高、RocketMQ 为什么在交易链路中更常见以及真实项目中的选型与优化思路。一、先看结论Kafka 和 RocketMQ 的差异本质上不是“谁更快”的问题而是两类设计目标的差异Kafka 把消息看作“不可变日志流”追求极致吞吐适合日志采集、埋点、流计算等场景。RocketMQ 把消息看作“关键业务事件”追求事务、延迟、顺序等业务语义适合订单、支付、库存等核心交易链路。如果只看纯吞吐6 节点集群下 Kafka 通常能做到 80 万到 120 万 TPSRocketMQ 大约是 30 万到 60 万 TPS相差 2 到 3 倍。但技术选型不能只看这一个指标还要看消息系统需要承载的业务约束。二、两种设计哲学Kafka只负责运货不关心内容Kafka 的核心能力来自日志模型每个 Partition 是一个有序、不可变、只追加写的日志段。消息写入就是顺序追加天然适合磁盘顺序 I/O。消费时通过sendfile把文件数据直接从内核送到网卡减少 CPU 拷贝。Producer 批量发送Consumer 顺序拉取链路设计高度统一。这套机制的前提是消息对 Kafka 而言是不需要理解的二进制流。它不需要解析业务字段也不需要关心消息之间的因果关系。RocketMQ为关键业务事件服务RocketMQ 需要让消息参与业务流程因此支持事务消息、延迟消息、顺序消息、Tag 过滤等能力。这些能力要求 Broker 在写入时理解消息内容因而无法完全复刻 Kafka 的纯日志路径。通俗地说Kafka 像高速收费站只做快速放行。RocketMQ 像全能服务区除了放行还要处理故障、加油、开票等业务动作。三、为什么 Kafka 吞吐更高1. 存储模型日志即消息 vs 日志加索引Kafka 的每个 Partition 对应独立日志目录消息的 offset 就是唯一标识写入路径很短。RocketMQ 使用CommitLog ConsumeQueue的两级结构CommitLog是全局共享日志所有 Topic 的消息都顺序写入这里。ConsumeQueue是每个队列的逻辑索引保存消息在CommitLog中的物理偏移量、消息大小和 Tag 哈希值。这意味着 RocketMQ 写入一条消息时除了写消息体还要异步或同步维护索引信息。它需要解析消息的 Topic、QueueId、Tag 等元数据因此消息必须先进入用户态内存。这也是 RocketMQ 写入路径无法完全零拷贝的架构原因。2. 写入路径纯追加 vs 用户态解析Kafka 的写入更接近“不关心内容”的纯追加网络接收 - 追加写日志 - PageCache - 异步刷盘RocketMQ 的写入路径更复杂网络接收 - 解析协议 - 路由到 Topic/Queue - 写 CommitLog - 构建 ConsumeQueue 索引 - 按需写 IndexFile/延迟索引其中“路由决策”和“索引构建”必须在用户态完成。Broker 要理解消息元数据才能决定它属于哪个逻辑队列以及如何构建后续查询索引。因此 RocketMQ 在写入路径上存在不可避免的 CPU 内存拷贝。3. I/O 优化顺序写、mmap 与 PageCacheKafka 充分利用了以下 I/O 特性顺序写磁盘避免随机寻道。使用mmap或sendfile减少内核态与用户态之间的数据复制。消息先进入操作系统 PageCache由系统决定刷盘时机同时降低 JVM GC 压力。RocketMQ 也使用顺序写和MappedByteBuffer但二级索引模型决定了它需要在用户态处理消息。即使消费路径可以使用FileChannel.transferTo()做零拷贝写入路径仍需要用户态解析和索引构建。4. 批量、压缩与传输链路Kafka 的 Producer 默认会攒批发送由batch.size和linger.ms控制。消息在批量后支持snappy、gzip、lz4等压缩可以减少网络传输和磁盘占用。消费端则通过sendfile实现零拷贝。RocketMQ 也支持批量发送和压缩但默认配置更保守。例如消息压缩需要显式配置compressMsgBodyOverHowmuch默认不开启。小消息高频发送时网络请求更多带宽利用率也更容易受影响。5. 水平扩展方式Kafka 的 Partition 是一等公民。Topic 创建后Partition 与 Broker 解耦后续可以通过kafka-reassign-partitions.sh将分区重新分配到新节点扩容相对平滑。RocketMQ 的 Queue 数量在 Topic 创建时确定并且会按当时存在的 Broker 进行分配。扩容 Broker 后已有 Topic 的旧 Queue 通常不会自动迁移新 Broker 主要承接之后创建的新 Topic。因此实践中往往会提前给高频 Topic 设置足够多的 Queue比如 64 或 128。四、为什么交易链路更常用 RocketMQ吞吐差距只在“疯狂写日志”这类纯吞吐场景下明显。进入电商交易场景后消息系统需要承载业务正确性RocketMQ 的强语义能力会变得更重要。1. 事务消息下单与扣库存一致下单场景需要保证“本地事务”和“消息发送”最终一致。RocketMQ 通过半消息和本地事务状态回查解决这个问题TransactionSendResultresultproducer.sendMessageInTransaction(msg,orderId);try{redisService.decrStock(orderId);returnLocalTransactionState.COMMIT_MESSAGE;}catch(Exceptione){returnLocalTransactionState.ROLLBACK_MESSAGE;}如果生产者宕机Broker 会定时回调CheckListener查询本地事务状态避免消息悬挂。Kafka 原生不提供这类事务消息模型虽然具备幂等 Producer 和 EOS但实现业务级事务通常需要 TCC、Saga 或补偿逻辑复杂度更高。2. 延迟消息订单超时关单“下单后 30 分钟未支付自动关单”是典型延迟投递场景。RocketMQ 内置 18 级延迟消息可以直接使用MessagemessagenewMessage(t_order_timeout,close_order,body);message.setDelayTimeLevel(7);producer.send(message);消息会先进入延迟队列到时间后再投递无需额外部署定时任务扫描数据库。Kafka 没有原生延迟队列能力通常需要外部调度器或时间轮方案补齐。3. 顺序消息订单状态不能乱序同一订单的状态流必须是“创建 - 支付 - 发货 - 完成”。RocketMQ 可以按订单 ID 把消息固定到同一个 Queue并使用顺序消费监听器MessageQueuemqmessageQueueList.get(orderId.hashCode()%queueNum);producer.send(message,mq);consumer.registerMessageListener(newMessageListenerOrderly(){OverridepublicConsumeOrderlyStatusconsumeMessage(ListMessageExtmsgs,ConsumeOrderlyContextcontext){returnConsumeOrderlyStatus.SUCCESS;}});Kafka 也可以按 Partition Key 做局部有序但 RocketMQ 的编程模型更贴近交易业务且支持广播、集群等更多消费模式。五、RocketMQ 的性能优化如果已经确定使用 RocketMQ可以从生产者、Broker、Topic 和集群四层做优化。1. 生产者批量与压缩DefaultMQProducerproducernewDefaultMQProducer(producer_group);producer.setSendMsgTimeout(3000);producer.setCompressMsgBodyOverHowmuch(1024);producer.getDefaultMQProducerImpl().setBatchMaxSize(512);producer.start();批量发送类似拼车发货单位成本更低。2. Broker异步刷盘与缓冲调大flushDiskType ASYNC_FLUSH writeBufferSize 64m mapedFileSizeCommitLog 1g异步刷盘适合日志、通知等允许短暂数据丢失风险的场景。对于资金类链路需要重新评估刷盘和复制策略。3. Topic扩充读写队列shmqadmin updateTopic-nlocalhost:9876-tt_gateway_log-r32-w32默认 8 个读写队列在高频场景下容易成为消费并发瓶颈可以按业务量扩到 32 或更高。4. 集群主从和多 Broker 部署通过主从 Broker 分摊写入与存储压力同时提升容灾能力。扩容时应结合 Topic Queue 规划避免“节点加了旧 Topic 没有负载迁移”的问题。5. 网关削峰大促零点流量通常不是均匀分布前几秒可能涌入大量订单请求。可以在 API 网关层加入基于 Disruptor 的环形缓冲把瞬时请求转换为批量推送RingBufferMessageEventringBufferRingBuffer.create(ProducerType.MULTI,MessageEvent::new,65536,newYieldingWaitStrategy());再配合定时批量发送和降级落库实现“宁可慢不可断”的削峰效果。六、选型建议场景类型推荐选型理由日志采集、埋点、流计算、数据管道Kafka极致吞吐I/O 链路高效订单、支付、库存、交易通知RocketMQ事务、延迟、顺序、Tag 过滤能力完整秒杀、大促、突发洪峰Kafka/RocketMQ 网关缓冲不依赖单点增加批量、缓冲与降级体系真正的技术选型不是比参数而是看消息系统在具体业务中要承担什么责任。Kafka 强在吞吐RocketMQ 强在语义表达和业务闭环。七、总结Kafka 的优势来自“数据即日志”顺序写、批量发送、零拷贝因此适合高吞吐数据管道。RocketMQ 的优势来自“事件即业务”通过事务消息、延迟消息、顺序消费和 Tag 过滤解决交易链路中的一致性与时序问题。它的吞吐低于 Kafka不是技术落后而是主动承担语义复杂性的结果。在大促等高并发场景中单靠消息中间件调优不够还需要在接入层增加缓冲、批量和降级能力。最终方案应服务业务本质日志管道用 Kafka交易中枢用 RocketMQ抗压体系靠系统级削峰。
返回列表