ARTICLE DETAIL

资讯详情

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

消息队列选型与零资损实践:RocketMQ如何扛住千万级实时互动

消息队列选型与零资损实践:RocketMQ如何扛住千万级实时互动 在母婴家庭赛道上“实时互动”和“零资损”往往是两件看似矛盾的事——家长上传一条宝宝视频希望瞬间推送到爷爷奶奶、外公外婆的手机上而一旦涉及会员购买、相册打印、保险这类交易又一条消息都不能丢。亲宝宝这个服务千万家庭的平台日均消息量早就不是小打小闹峰值时段甚至能冲到每秒几十万条。团队在选型时反复纠结于 Kafka、RabbitMQ、RocketMQ最终落定 RocketMQ 作为核心消息中间件把“秒触达”和“零资损”这两件事都扛住了。这篇文章不聊虚的直接把我整理的消息队列选型对比、消息链路设计、事务消息保障、以及部署调优过程中踩过的坑都摊开讲给正在做同类业务场景的工程团队一些可以直接抄作业的参考。1. 从“选型纠结”到“一锤定音”RocketMQ 是如何在 Kafka 和 RabbitMQ 中胜出的1.1 千万家庭场景到底对消息中间件提了什么要求先还原一下业务场景。亲宝宝的核心玩法是“家庭共同育儿”家长创建家庭组把爷爷奶奶、外公外婆、亲戚朋友拉进来日常上传宝宝照片、视频、身高体重记录、成长日记其他家庭成员可以点赞、评论、互动。表面上看这像一个社交 Feed 场景但背后还挂着付费会员、照片打印、母婴保险等高价值交易业务。也就是说消息中间件一条链路要同时服务两类极端需求实时互动链路点赞、评论、新内容通知必须秒级触达用户端延迟高了家长会明显感知“手机没响”。资金/权益链路订单支付成功、权益发放、发货通知任何一条消息丢失或重复都可能造成真金白银的资损。再加上千万级家庭意味着亿级用户图片视频元数据、关注关系、动态流、行为事件这些消息类型相互交织峰值流量呈脉冲式爆发。工程团队对消息中间件的核心诉求可以归纳为四点低延迟、高可靠、可追踪、易运维。这四点看单都不难但叠加在一起选型就没那么简单了。1.2 三个主流消息队列的横向对比我以实战视角把 Kafka、RabbitMQ、RocketMQ 放在一起对比直接给出结论对比维度KafkaRabbitMQRocketMQ吞吐能力最强单集群轻松百万级 TPS中等万级 TPS 后瓶颈明显较强十万级 TPS 很稳消息延迟正常几十毫秒峰值下易波动微秒级到毫秒级但堆积后劣化毫秒级长轮询推送稳定事务消息事务 API 存在跨分区一致性实现复杂弱依赖插件或上层补偿原生支持半消息 状态回查落地友好消息轨迹需要额外采集插件支持能力有限内置轨迹链路开箱即用死信队列有 Dead Letter Queue但消费重试策略简单有 DLX功能完整有重试队列 死信队列配置灵活运维复杂度依赖 Zookeeper组件多轻量但镜像队列集群坑多天生无 ZookeeperNameServer Broker 架构简单我们这个场景最敏感的是“零资损”。Kafka 虽然吞吐最强但它的“至少一次”投递语义加上分区级别的有序性导致消费端写业务库时天然要面对大量乱序和重复事务机制用起来心智负担很重。RabbitMQ 的路由模型确实灵活但想要支撑亿级用户的消息洪峰集群扩展和消息堆积能力很快会成为短板。RocketMQ 的好处在于它把消息中间件从“管道”做成了“平台”事务消息、定时消息、消息轨迹、消费重试、死信队列这些能力都内置不用再自己造轮子补丁式地保证不丢不重。1.3 为什么事务消息成为决定性砝码线上交易场景最怕什么最怕“本地数据库事务提交了消息却没发出去”或者反过来“消息发出去了本地事务回滚了”。这两者都会造成业务状态不一致进而导致资损。RabbitMQ 有 publisher confirm 但缺少实用的分布式事务方案Kafka 的事务 API 需要配合 Kafka Streams 使用而且语义较重。RocketMQ 的“半消息”机制则把分布式事务拆成了发送预消息、执行本地事务、提交/回滚、失败回查四步业务代码可以用相对优雅的方式保证订单表和消息表最终一致。我在调研时的判断是如果只做日志采集、行为上报Kafka 无疑是最优解如果只做内部系统解耦RabbitMQ 也够用但既要实时互动又要交易一致性RocketMQ 是国内团队最顺手的选择。后来也印证了这一点RocketMQ 在阿里内部有双十一那么极端的流量验证社区活跃度又高排查问题时搜得到大量中文踩坑经验这一点在深夜值班时尤其值钱。2. 实时互动的消息链路设计从 Topic 划分到推送触达2.1 消息分类与 Topic 规划选型定下来之后第一步不是急着写代码而是做消息域划分。我把亲宝宝场景里的消息分成三类行为事件消息点赞、评论、收藏、关注、浏览记录这类消息量大、容忍一定延迟主要用于计算 Feed 流和推送通知。内容同步消息宝宝照片、视频、成长记录上传后需要同步给相册服务、ES 索引、推荐服务、CDN 刷新服务这类消息要求可靠不能丢。交易与权益消息订单支付、会员开通、保险投保、发货状态变更这类消息必须做到不丢不重且要保留完整轨迹用于审计。对应的 Topic 合理切分是Topic 名称消息类型Tag 示例主要消费者family_event行为事件like、comment、follow、share通知服务、Feed 服务growth_record内容同步photo_uploaded、video_uploaded、record_created索引服务、CDN 服务order_event交易消息order_paid、membership_activated、shipment_status订单服务、权益服务、对账服务这里一个关键经验Topic 数量不要贪多按主业务域划分即可细粒度用 Tag 解决。RocketMQ 的 Topic 底层是读写队列Topic 太多会稀释队列资源管理成本和 Broker 压力都会上升。Tag 用于消费端过滤但要注意它是在客户端拉取时过滤不是完全省流量的方案真正过滤源头项是合理拆 Topic。2.2 消息体设计统一规范消息体必须标准化否则一百个服务各发各的消费端没法处理。我给出的推荐格式是一个扁平 JSON包含基础字段和业务 payload{ eventId: uuid-随机全局唯一, bizId: 订单号或内容ID, bizType: order_paid, userId: 用户ID, familyId: 家庭组ID, timestamp: 1710000000000, source: order-service, payload: { orderId: M202501010001, amount: 19900, productType: photo_print } }这里的eventId全链路唯一用于消息轨迹追踪bizId是业务幂等键消费端去重靠它timestamp最好用发送方时间戳否则跨服传输后消费端拿到的时间可能失真。我们内部把这套格式称作“消息公约”每个服务发消息前都要校验 JSON 结构不符合规范直接报错从源头避免下游解析崩溃。2.3 生产端发送策略秒触达的第一公里实时互动的“秒”主要比拼的是生产端到消费端全链路延迟。生产端常用的是DefaultMQProducer发送方式有同步、异步、单向三种。对于实时通知类消息我用同步发送但设置了 3000ms 超时兜底异步发送虽然性能更好但回调链路容易造成线程阻塞和历史问题难追查我在交易场景不推荐。DefaultMQProducer producer new DefaultMQProducer(order-event-producer); producer.setNamesrvAddr(namesrv1:9876;namesrv2:9876); producer.setSendMsgTimeout(3000); producer.setRetryTimesWhenSendFailed(3); producer.start(); Message message new Message(order_event, order_paid, body); SendResult result producer.send(message); if (result.getSendStatus() SendStatus.SEND_OK) { // 记录发送成功日志 }注意setRetryTimesWhenSendFailed这个参数RocketMQ 默认失败会重试 2 次但重试可能导致重复发送。配合消费端幂等这里可以放心开重试总比重试次数过少导致丢消息强。2.4 消费端水平扩展与长轮询机制RocketMQ 的 PushConsumer 之所以能实现“秒触达”不是靠 broker 主动 TCP 推送而是靠长轮询客户端向 broker 拉消息时如果队列为空会挂起一段时间等待数据到达一旦有消息写入立即可返回。这个机制让 broker 的推送延迟实际可控在几十毫秒级。消费端的水平扩容有一个约束RocketMQ 的消费模型是队列粒度负载均衡一个消费者组内的实例数不能超过订阅主题的队列总数否则多出来的实例会空转。比如某个 Topic 队列数配置为 16消费者实例最多起 16 个想继续扩就要先把 Topic 的读写队列数调大。我建议核心实时链路的 Topic 队列数至少 32 以上给双十一这类峰值场景预留扩展空间。还需要调优的消费端参数包括consumeThreadMin/consumeThreadMax消费线程数实时场景设为 16~32交易场景建议 8~16防止线程过多压垮下游数据库。consumeMessageBatchMaxSize批量消费条数默认 1实时场景调到 8~32 提升吞吐但要注意批量内某条失败会影响整批重试。maxReconsumeTimes消费重试次数默认 16 次交易场景建议设 3~5 次避免无效消息占用队列资源。实操中我们发现真正让端到端延迟从“秒”变成“百毫秒”的关键不在 RocketMQ 本身而在消费端代码是否轻量。点赞、评论这类消息的处理逻辑如果做成同步调用通知网关、再通知网关调 Android/iOS PUSH链路一长延迟就会被放大。正确做法是消费端只做核心状态变更和落库然后立刻 ACK再通过另一条异步任务或 RocketMQ 的延迟消息去驱动 PUSH 任务。3. 零资损如何落地事务消息、幂等消费和定时对账的三层防线3.1 事务消息的完整生命周期资损防线的第一层是发送端的事务一致性。以用户购买会员为例支付回调后订单服务需要做两件事修改订单状态为“已支付”发消息通知权益服务开通会员。如果先改库再发消息消息发送失败会导致权益未开通如果先发消息再改库消费者看到消息查询订单却是未支付就会错乱。RocketMQ 事务消息的流程可以归纳为几步生产者向 Broker 发送“半消息”Half Message此时消息对消费者不可见。发送成功后生产者执行本地事务例如更新订单表状态、写本地事务日志表。本地事务成功则提交commitBroker 让消息可见失败则回滚rollbackBroker 丢弃半消息。如果提交/回滚动作因为网络原因没传到 BrokerBroker 会在一段时间后向生产者发起回查check back生产者根据本地事务表状态决定提交还是回滚。实践时我会额外维护一张transaction_log表记录事务记录 id、事件类型、业务主键 id、状态。回查逻辑查这张表即可不需要去业务表里猜。public class OrderTransactionListener implements TransactionListener { Override public LocalTransactionState executeLocalTransaction(Message message, Object arg) { // 在这里执行本地事务更新订单状态 写入 transaction_log try { orderService.markPaid(orderId); transactionLogService.writeSuccess(orderId, message.getKeys()); return LocalTransactionState.COMMIT_MESSAGE; } catch (Exception e) { transactionLogService.writeUnknown(orderId, message.getKeys()); return LocalTransactionState.UNKNOW; } } Override public LocalTransactionState checkLocalTransaction(MessageExt message) { // 回查逻辑查询 transaction_log 是否已成功 Long orderId Long.parseLong(message.getKeys()); TransactionStatus status transactionLogService.getStatus(orderId); if (status TransactionStatus.SUCCESS) { return LocalTransactionState.COMMIT_MESSAGE; } return LocalTransactionState.ROLLBACK_MESSAGE; } }事务消息能解决“本地事务和消息发送的原子性”但它不是银弹。生产上还要为回查接口设置超时和降级防止 Broker 回查风暴把事务库打爆。回查间隔和最大回查次数都可以在 broker 配置中调我默认保持 60s 一次回查次数不要超过 20 次否则异常消息长期占着半消息队列会拖垮 broker 资源。3.2 消费端幂等重复消息必须天然可容忍即使发送端做到事务一致消息从 Broker 到消费者还是会因为网络重试、消费者宕机恢复、Rebalance 产生重复投递。RocketMQ 默认是“至少一次”语义所以消费端必须幂等。我不建议把幂等完全押在 RocketMQ 自带的去重机制上而是业务侧用唯一键约束。最稳妥的做法是在消费者本地写一张consume_record表将bizType bizId作为唯一索引。消费时先插入记录插入成功才执行业务逻辑插入失败立刻 ACK当作已处理。这个方案的好处是无论消息重复多少次数据库唯一索引都会挡住第二次。CREATE TABLE consume_record ( id bigint(20) NOT NULL AUTO_INCREMENT, biz_type varchar(32) NOT NULL COMMENT 业务类型, biz_id varchar(64) NOT NULL COMMENT 业务幂等ID, message_id varchar(64) NOT NULL COMMENT RocketMQ消息ID, status tinyint(4) NOT NULL DEFAULT 1 COMMENT 1已消费, create_time datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (id), UNIQUE KEY uk_biz_type_biz_id (biz_type,biz_id) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;如果不想引入表也可以用 RedisSET NX做短时间幂等比如设置 24 小时过期。不过 Redis 方案在极端情况下有缓存淘汰风险涉及资金权益的业务我坚持走数据库唯一索引。3.3 消费重试与死信队列的告警闭环RocketMQ 的消费失败默认会按照延迟级别重试 16 次但重试策略必须结合业务场景调。比如order_paid消息消费失败可能是因为下游权益服务暂时抖动重试几次就能好但如果是因为消息体本身有脏数据重试到死也不会成功还会白白消耗 broker 吞吐。我的做法是分场景定制基础通知类消息重试 3 次失败进死信队列。交易核心消息重试 5 次失败进死信队列同时立刻发送告警。对账/索引同步类消息重试 16 次也不怕因为这类任务允许长延迟。死信队列的消息不能被默认消费者继续消费否则会无限循环。RocketMQ 控制台可以查看死信消息并重置消费位点但线上建议直接消费死信主题解析原始消息后调用补偿接口进行处理。我们的补偿服务专门消费%DLQ%order_event这类主题重放前先查consume_record判断是否已经处理过避免人工重放产生二次资损。3.4 定时对账最后一道保险再完美的技术方案也可能有极端情况比如 broker 磁盘损坏、人工误操作删了 Topic。所以必须要有定时对账。我们内部跑一个每日任务扫描发送端transaction_log表把所有状态为已完成但对应业务表没有预期变化的数据捞出来通过重新发送消息或者调用补偿接口修复。此外还有一条链路rocketmq 的消息轨迹数据从trace主题导出到数仓然后和订单表做关联比对。通过轨迹能知道消息何时发送、何时被谁消费、消费是否成功一旦发现某个订单支付消息没有对应的消费 ACK立刻触发补偿。资损问题不是靠单一技术解决的事务消息解决“发不发”一致性的问题幂等消费解决“重不重”的问题定时对账解决“万一还是漏了”的问题三层防线全部落地才算真正敢说一句“零资损”。4. 部署、调优、排坑从环境准备到高可用集群的实战经验4.1 CentOS 7 环境下的安装部署要点尽管标题押在业务上但支撑千万家庭互动的前提是集群本身稳定可靠。很多团队第一步就栽在安装和部署上我先给出一个从零搭建的最小可运行方案。环境准备建议JDK 1.8Linux 内核 3.10 以上的 CentOS 7。下载 RocketMQ 时注意版本4.9.x 和 5.x 的配置文件格式差异不小新手建议先跑通 4.9.x再看 5.x 的 Controller 模式。解压后需要修改两处脚本# runserver.sh 中 NameServer 的 JVM 内存建议改为 2g JAVA_OPT${JAVA_OPT} -server -Xms2g -Xmx2g -Xmn1g # runbroker.sh 中 Broker 的 JVM 内存建议改为 8g 以内 JAVA_OPT${JAVA_OPT} -server -Xms4g -Xmx4g -Xmn2g为什么 JVM 内存不能一上来就调很大RocketMQ 的大量读写依赖 PageCacheJVM 堆内存占太多会挤掉 PageCache 的空间反而导致刷盘变慢。堆内存 8G、PageCache 预留 16G 以上是比较健康的比例。启动顺序是先启动 NameServer再用nohup sh bin/mqbroker -c conf/broker.conf启动 Broker。注意 broker.conf 要显式配置brokerClusterNameDefaultCluster brokerNamebroker-a brokerId0 brokerIP110.0.0.11 listenPort10911 deleteWhen04 fileReservedTime48 brokerRoleSYNC_MASTER flushDiskTypeSYNC_FLUSHbrokerIP1必须显式指定这是非常多人踩过的坑。如果服务器有多个网卡或部署在 Docker 里RocketMQ 启动时会自动探测 IP很容易注册成内网容器 IP 或公网 IP导致客户端连不上、消息发送超时。我们就在这个brokerIP1上吃过亏换了一个网卡环境后客户端全部报connect to xx timed out排查半天才发现 broker 向 NameServer 注册的 IP 是错的。4.2 高可用架构同步复制还是 Dledger如果业务不能容忍消息丢失broker 主从架构必须开同步复制。传统的 Master-Slave 模式有 ASYNC_MASTER 和 SYNC_MASTER 两种。ASYNC_MASTER 模式下消息先落到主节点就返回成功从节点异步复制主节点宕机且从节点还没来得及同步时丢失的是这一小段消息。SYNC_MASTER 模式下消息会等待从节点也写入成功后返回吞吐量会有所下降但几乎不丢数据。对于“零资损”目标我强烈建议brokerRoleSYNC_MASTERflushDiskTypeSYNC_FLUSH双保险。不过传统主从还有一个痛点主节点宕机后从节点只读不可写需要人工切换。如果想要自动故障转移就上 DledgerRaft 模式官方 RocketMQ 4.5 开始支持至少 3 个节点选主自动实现。代价是部署更复杂机器成本更高。我们的容忍度是核心交易集群用 Dledger而实时互动集群用 SYNC_MASTER 主从兼顾可靠性和成本。4.3 高频踩坑和排查链路把我在实战中遇到的几个高频问题列出来每条都是真金白银换来的消息堆积排查思路先看消费者在线数量、消费 TPS 和堆积量。如果消费者数量小于队列数优先扩容实例如果消费 TPS 很低说明消费逻辑有瓶颈需要看下游 SQL 慢日志、Redis 阻塞、锁竞争等问题。RocketMQ 控制台提供了详尽的消费进度会显示每个队列的diffTotal按这个数字排序就能快速锁定异常 Topic。Rebalance 导致的重复消费消费者实例变动时队列的所有权会重新分配。此时可能上一实例还没处理完的消息被另一实例重新拉取重复消费不可避免。这就是为什么我一直强调消费端一定要做幂等而不是心存侥幸。消息乱序RocketMQ 只保证队列级别的有序不同队列之间不保证。要实现同一个业务 ID比如同一个订单的消息有序生产端必须用MessageQueueSelector把相同 ID 的消息都发到同一个队列。我们做会员续费通知时就遇到过用户同时在 App 和 H5 端触发支付两条消息被分到不同队列消费端先后顺序错乱导致权益开了又关的问题。后来改造为以userId作为选择 key才从根上解决。多网卡 IP 注册问题前面提到过brokerIP1这里再补充一种情况RocketMQ 4.x 版本在 Docker 容器中默认brokerIP1是容器 IP需要在环境变量ROCKETMQ_BROKER_IP指定宿主机 IP否则外部客户端无法访问 broker。排查这类问题时在客户端机器执行telnet brokerIP 10911是最直观的验证方式。4.4 压测与容量预估上线前必须压测。我用 JMeter rocketmq-client 压测过的经验来看核心指标就三个生产端发送 TPS、消费端消费 TPS、端到端延迟 P99。建议分两个阶段单 Topic 单队列压测摸清 broker 单节点性能基线。多 Topic 多队列混合压测模拟真实业务混合流量。压测时关注 broker 的commitlog写入性能和 PageCache 命中率如果机器内存不够就会频繁刷盘延迟曲线明显飙升。容量预估公式我习惯用“峰值 TPS × 消息平均大小 磁盘吞吐”再预留 30% 冗余。比如峰值 10 万 TPS、单条 1KB那么每秒要写入 100MB选 SSD 并且硬盘 IOPS 至少 10000否则消费端还没积压broker 先扛不住了。5. 从运营视角看实时互动与消息可靠性的平衡技术方案落地之后还要面对一个长期问题实时互动和消息可靠性在某些场景下其实是互斥的。比如为了保障交易消息不丢我们开了同步复制和同步刷盘但同样的配置如果应用到点赞、评论消息上成本会翻倍。所以集群要按业务分级治理。我们做了两套 RocketMQ 集群的逻辑隔离实时互动集群异步刷盘 主从异步复制保证低延迟能容忍极少量的消息丢失。资金交易集群同步刷盘 同步复制宁可吞吐低一点也要保证不丢消息。Topic 的分配也有原则涉及资金变动、权益发放的消息只允许发往交易集群行为事件消息只允许发往互动集群。通过生产端客户端命名空间隔离和 broker 的rejectPublicly策略从接入层面禁止串用防止某个开发同学图省事把订单消息发到异步集群。除此之外实时互动场景还有个容易被忽略的点家长给宝宝的成长记录点赞、评论这种消息延迟 1 秒还是 5 秒用户体感差异不大但如果你在点赞消息里挂了“发送微信服务号通知”这种外部依赖那延迟就会被外部服务的不稳定性放大。我后来把PUSH通知从主消费链路中拆出来放到 RocketMQ 的延迟消息级别延迟个 30 秒再触发推送既保证了用户收到通知又不影响主链路 ACK 速度。这种“降级式”设计比单纯调大消费线程数靠谱得多。最后再分享一个压箱底的小技巧消费端代码里不要直接消费消息对象然后裸写业务代码而是把消息先转换成一个“领域事件”对象所有消费者统一走一个事件分发器。这样后续如果从 RocketMQ 切换到其他中间件或者一份消息要同时触发多个处理器只需要改分发器不需要动业务逻辑。我们在亲宝宝的互动链路上就是这么设计的几次业务迭代和中间件参数调整业务代码基本没有大改。消息中间件终究只是手段真正决定系统能不能扛住千万级流量和“零资损”要求的是围绕它建立的一整套路由、事务、幂等、对账体系。把这套体系梳理清楚RocketMQ 才能从“一个消息队列”变成业务增长的稳定底座。
返回列表