ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

RocketMQ消息轨迹追踪:埋点链路、异步上报与生产实践

RocketMQ消息轨迹追踪:埋点链路、异步上报与生产实践 一次Java后端面试里聊完RocketMQ的基础架构和顺序消息后面试官顺势抛出一句“那你说说RocketMQ的消息轨迹追踪是怎么实现的”这个问题表面上是在问一个功能点实际上考的是你对RocketMQ客户端埋点机制、异步上报链路和工程取舍的理解深度。很多候选人能说出“默认Topic叫RMQ_SYS_TRACE_TOPIC”这种话但一旦被追问“轨迹数据是在哪个环节产生的为什么会丢对性能影响怎么评估”就卡住了。这篇文章就围绕这个问题把RocketMQ消息轨迹从设计动机、核心数据结构、源码链路到生产配置完整拆一遍。如果你正在准备Java中间件方向的面试或者你在实际项目里被“某条消息到底有没有被消费过”折磨过这篇应该能帮你把脑子里的知识点串成一条线。1. 面试官真正想听的消息轨迹到底解决了什么问题1.1 没有轨迹时排查一条消息要费多大劲我在实际维护消息中间件相关的服务时最怕的不是消息积压而是“消息好像丢了”。业务流程是订单支付成功后发送一条延迟消息用来通知库存系统去锁库存结果库存那边一直没反应。这时候你要定位问题只能翻三份东西生产者应用日志、Broker的存储日志、消费端应用日志。三份日志分布在不同的机器上靠消息里的业务Key一个个串耗时不说还很容易被“消息在什么时候进入消费端、消费失败后又重试了几次”这种跨端信息卡住。如果打开了消息轨迹情况完全不一样。你能直接看到这条消息从生产者发出、写入Broker、被推给消费者、消费者返回成功或失败的时间线而且这些信息集中在一个地方按Message ID或者Key一搜就出来。面试官问这个问题潜意识里的考察点其实是你有没有真的用消息中间件解决过线上的诡异问题还是只停留在写写Producer、Consumer的demo阶段。1.2 消息轨迹和全链路追踪不是一回事很多人会把RocketMQ的消息轨迹和平时用的链路追踪类似基于TraceId的分布式调用链混在一起这是面试里很常见的扣分点。链路追踪关注的是“一次业务请求跨了哪些服务”它记录的是RPC调用级别的Span而消息轨迹关注的是“一条消息在MQ内部和客户端两侧的关键事件”它至少要覆盖发送前、发送后、消费前、消费后四个节点。你可以把消息轨迹理解成快递物流一个包裹发出去了中途在每个中转站扫描一次最终签收或者拒收每一步都有时间戳和状态。链路追踪则是你查“这个包裹对应的订单是怎么一路流转过来的”。两者可以配合但不能互相替代。RocketMQ的消息轨迹更偏向“包裹物流”它不管业务系统内部的调用关系只管消息本身的生命周期。1.3 面试考核的三层能力我把这个问题拆成三个层次层次考察内容面试官想听的回答第一层知不知道有这个功能“RocketMQ提供了消息轨迹默认Topic是RMQ_SYS_TRACE_TOPIC可以追踪消息发送和消费的状态。”第二层知不知道实现原理“客户端通过Hook在发送/消费前后埋点把轨迹数据异步批量上报到轨迹Topic再由控制台或者自定义程序消费展示。”第三层知不知道工程取舍“轨迹是诊断数据不是业务数据丢了不影响主流程异步批量上报是为了减少性能损耗有界队列满了会丢弃轨迹而不是阻塞业务。”能说到第三层基本就过了。2. 消息轨迹的实现骨架三个角色和一条异步流水线2.1 核心角色Hook、TraceContext、TraceDispatcherRocketMQ消息轨迹的实现不复杂核心就三个部分第一部分是埋点Hook。客户端定义了一组钩子接口比如SendMessageTraceHookImpl负责在生产者发送消息前后被回调ConsumeMessageTraceHookImpl负责在消费者消费前后被回调。这其实和Servlet的Filter、Spring的AOP是同一个思路不侵入你的业务代码而是在框架的固定节点插入“旁路逻辑”。第二部分是轨迹上下文。每次埋点会生成一个TraceContext对象里面装的是这次发送或消费的现场快照消息ID、Topic、Group、客户端地址、时间戳、耗时、成功失败状态等。它内部又包含了一个TraceBean列表每个TraceBean对应一条被追踪的消息。第三部分是异步分发器。TraceContext不会立刻被发出去而是被丢进一个AsyncTraceDispatcher内部的队列。Dispatcher里有一个后台线程攒够一批数据或者每隔几秒就批量发送一次把轨迹数据写到RocketMQ自己的轨迹Topic里。这样就把轨迹上报的耗时和业务路径解耦了。这三者的关系可以类比成快递员Hook在揽件时填单子TraceContext然后把单子统一放进中转站的仓库队列再由一辆定时发车的货车Dispatcher批量运走。2.2 TraceContext里到底装了什么数据面试的时候如果能随口说出几个关键字段会显得你是真的看过源码。我整理了一张常见字段表字段含义traceType轨迹类型Pub表示发送、SubBefore表示消费前、SubAfter表示消费后groupName生产者或消费者所在的Grouptopic消息所属的业务TopicmsgId客户端的消息IDoffsetMsgId消息在Broker上的物理偏移ID通常查询时更可靠clientHost客户端IP和端口bornTime消息在生产端创建的时间storeTime消息在Broker落盘的时间costTime该阶段的耗时比如发送耗时或消费耗时isSuccess当前阶段是否成功面试官如果追问“为什么一个TraceContext里会有TraceBean列表而不是单个Bean”你可以解释批量消息发送或者批量消费时一次动作可能涉及多条消息所以一个上下文对应多个TraceBean。这种数据模型的设计是为了覆盖批量场景。2.3 为什么上报必须是异步且批量这个问题几乎是必问的。轨迹上报本质上是“额外写一份日志”如果每个消息发送后都同步再发一条轨迹消息等于让原本一次消息发送变成两次网络交互TPS直接腰斩不说还会让业务线程卡在轨迹发送上。所以RocketMQ做了一定程度的妥协轨迹数据进入本地有界队列后立即返回业务线程不等待后台线程按批聚合攒够一定数量比如100条或者到达时间阈值再统一发送轨迹发送失败只记录WARN日志不回滚、不重试或者有限重试因为它丢了不影响业务。这就引出一个面试加分项有界队列满了怎么办答案是丢弃新进来的轨迹数据同时打印告警日志。这样做的逻辑是轨迹数据是诊断数据丢掉几条最多让你在排查问题时少点现场信息但如果你因为轨迹上报把业务线程阻塞了那才是真正的生产事故。这种“非核心数据可丢失”的设计思想在很多中间件里都有体现。3. 跟读源码从SendMessageTraceHook到TraceDispatcher的完整路径3.1 发送一条普通消息时钩子做了什么如果你打开RocketMQ客户端源码在org.apache.rocketmq.client.trace.hook包下能看到两个关键实现类SendMessageTraceHookImpl和ConsumeMessageTraceHookImpl。这两个类分别实现SendMessageHook和ConsumeMessageHook接口接口里都有before和after两个方法。拿发送场景举例。生产者的send方法在真正网络发送前会先调用注册进来的SendMessageTraceHookImpl.beforeSendMessage这个阶段会做什么它会把当前时间、客户端地址、消息ID、Topic这些信息先填充到一个TraceContext里此时costTime还是0状态也是未知。之后消息真正发送Broker返回发送结果这时钩子的afterSendMessage被调用它会把发送结果成功还是失败、耗时算出来再把这个TraceContext交给AsyncTraceDispatcher。这里有个细节after阶段拿到的offsetMsgId是Broker返回的物理偏移ID比msgId更能定位消息在Broker上的真实位置所以在追踪数据里两个ID最好都存。消费端也是类似的逻辑。消费者在回调你的MessageListener之前会先触发beforeConsumeMessage生成一个traceTypeSubBefore的上下文等你的监听器返回CONSUME_SUCCESS或RECONSUME_LATER之后会触发afterConsumeMessage生成一个traceTypeSubAfter的上下文里面记下消费状态和耗时。所以一条消息如果消费失败并重试了三次你在轨迹里会看到一条发送轨迹外加三条“消费前”和三条“消费后”的轨迹记录。3.2 AsyncTraceDispatcher内部的攒批逻辑我最初看这块代码时以为轨迹数据是来一条发一条后来才发现自己格局小了。AsyncTraceDispatcher内部维护了一个有界队列所有Hook产生的TraceContext都会被append到这个队列里。Dispatcher里有个后台线程flushRunnable它的工作很简单从队列里批量拉取TraceContext按上下文里的信息组装成满足批量大小的消息列表把这一批消息包装成MessageTopic是指定的轨迹Topic默认RMQ_SYS_TRACE_TOPIC通过一个内部的生产者实例发送出去。这个后台线程有两个触发条件一个是攒够批量大小不同版本默认值略有差异通常默认100条左右一个是到达固定的flush间隔通常几秒钟。谁先满足谁触发。这样的设计保证了两个极端场景高吞吐下轨迹不会积压太多低吞吐下轨迹不会由于一直攒不够数量而迟迟不落库。需要特别注意的是这个内部生产者发送轨迹消息时本身也会走SendMessageHook。如果不加控制就会造成“轨迹消息的轨迹消息”这种递归埋点。我没细看源码时也踩过这个认知坑实际上实现里会把内部生产者的钩子关闭或者标记为不用再追踪避免循环上报。3.3 轨迹数据本身也是普通消息很多面试者会把消息轨迹想得很神秘其实轨迹数据落到Broker之后跟一条普通消息没有任何区别。它照样写入CommitLog照样被复制到从节点照样按Topic维度去消费。你甚至可以直接用一个Consumer去订阅RMQ_SYS_TRACE_TOPIC把轨迹数据通过日志或者数据库自己沉淀下来。官方控制台的消息轨迹查询本质上就是消费这个Topic里的轨迹消息然后按Message ID或消息Key把同一批轨迹串出来渲染成时间线。弄懂这一点很多问题就顺了为什么轨迹数据默认只保留几天因为RocketMQ对消息的清理是按CommitLog文件保留时间或磁盘水位来做的轨迹Topic并没有特殊的“延长保留”待遇。为什么说开启轨迹会增加磁盘开销因为轨迹消息也是消息也会占存储。这些我在后面的实操部分再展开。4. 落地实操开启轨迹、查询轨迹与验证实验4.1 开启轨迹前需要想清楚的三件事第一你的磁盘扛不扛得住。轨迹消息虽然小但高频业务下数量惊人。假设你的业务Topic每秒产生1万条消息开启轨迹后相当于系统里每秒又多了1万条小消息。如果这些轨迹消息都落在同一块磁盘上会显著增加磁盘IO压力。第二选默认轨迹Topic还是自定义Topic。默认RMQ_SYS_TRACE_TOPIC的优点是配置简单、控制台开箱即用缺点是所有开启轨迹的业务都往同一个Topic里写数据混在一起权限也不好隔离。自定义轨迹Topic的优点是你可以按业务拆分、方便单独制定清理策略缺点是要多建一个Topic。第三判断哪些客户端需要开启。很多生产事故其实是“生产端开了轨迹消费端没开”或者反一反。这样轨迹只能看一半等于白开。开启时要确保一条业务链路涉及的生产者和消费者都打开才能形成完整时间线。4.2 具体配置步骤RocketMQ消息轨迹的开关不只是客户端的事Broker侧也有配置。最基础的配置是在broker.conf里加一行traceTopicEnabletrue配置完需要重启Broker生效或者在已经启动的Broker上通过更新配置的方式动态打开不同版本的处理方式不太一样建议以你所用版本官方文档为准。某些版本还需要在NameServer侧也配置traceTopicEnabletrue否则客户端在拉取轨迹Topic路由时会遇到问题。客户端侧生产者的开启方式有两种等价写法。第一种是在构造方法里直接指定DefaultMQProducer producer new DefaultMQProducer( order_pay_group, true, // enableMsgTrace order_pay_trace_topic // 自定义轨迹Topic可空 ); producer.setNamesrvAddr(127.0.0.1:9876); producer.start();第二种是默认构造对象以后通过setter打开DefaultMQProducer producer new DefaultMQProducer(order_pay_group); producer.setNamesrvAddr(127.0.0.1:9876); producer.setEnableMsgTrace(true); producer.setCustomizedTraceTopic(order_pay_trace_topic); producer.start();消费者端同样支持DefaultMQPushConsumer consumer new DefaultMQPushConsumer( order_pay_consumer_group, true, order_pay_trace_topic ); consumer.setNamesrvAddr(127.0.0.1:9876); consumer.subscribe(order_pay_topic, *); consumer.setMessageListener((msgs, context) - { // 业务处理 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; }); consumer.start();如果自定义了轨迹Topic需要提前确认这个Topic在Broker上已经创建否则内部生产者发消息时会因为Topic路由不存在而报错。4.3 控制台查询轨迹看什么、怎么看在RocketMQ Dashboard里找到“消息轨迹”或者“Trace”查询入口填入你关心的Message ID或者消息Key就能看到这条消息的轨迹时间线。我一般会按以下几个关键信息去判断问题先看发送轨迹里有没有costTime异常大的记录发送耗时过大通常说明客户端到Broker的网络链路有问题或者Broker端写盘慢。再看SubBefore和SubAfter的次数如果一条消息有多次消费轨迹说明消费端返回过RECONSUME_LATER经历过重试。最后看SubAfter的isSuccess如果消费最终失败时间线里会看到最后一次失败的时间点配合业务日志里的异常栈去定位根本原因。我记得有一次排查延迟消息未触发的场景生产端显示消息已经发送成功消费端始终没有SubBefore记录。最终发现是消费者的线程池被某个慢任务占满消息一直在Broker端排队等待投递但轨迹里看不到“排队中”的状态需要你结合消费组消费并发度和积压数去反推。那条消息的轨迹价值在于帮我确认了“不是消息丢了而是消费端自己堵住了”。4.4 自己搭一个最小验证实验如果你手边有RocketMQ环境我建议花半小时做一个小实验写一个正常消费者、一个故意抛异常的消费者分别订阅同一个Topic然后往Topic里发两条消息。正常消费者的轨迹你会看到发送成功一条、消费前一条、消费后一条消费后的isSuccess为true。异常消费者那条消息的轨迹则是发送成功后会出现多条SubBefore和SubAfter记录SubAfter的isSuccess为false直到重试达到上限后进入死信队列。通过这个对比实验再回头去看TraceType枚举那三个值印象会深刻得多。5. 追问环节轨迹的边界、代价和工程取舍5.1 追问一开启消息轨迹后性能损耗到底有多少我见过的面试回答大多是“有损耗但不大”这种回答太糊了。更好的说法是分点拆开第一客户端侧损耗是“每个消息多一次对象创建和时间戳记录”这笔开销很小第二真正成本在异步传输和Broker存储因为轨迹消息会增加网络包数量和CommitLog写入量第三由于批量上报客户端侧不会出现“一条业务消息对应一次额外网络请求”的放大效应。如果你的业务消息体本身很小、TPS又特别高轨迹消息带来的额外写入比例会被放大。比如业务消息只有几百字节而轨迹消息可能也有几百字节相当于存储写入翻倍这种场景下就要考虑只对核心链路开启轨迹或者通过自定义Hook实现采样上报只记录一定比例的轨迹。5.2 追问二轨迹数据丢了怎么办面试官问这个问题其实是想看你有没有分清业务数据和诊断数据。最好这样答轨迹数据的定位是“诊断辅助”它允许丢、允许不完整真正要保证不丢的是业务消息本身。所以当轨迹上报队列满了或者上报失败的时候RocketMQ选择丢弃并告警而不是阻塞业务发送。如果你有强审计需求比如金融场景里要求每条消息都留痕那就不能只依赖RocketMQ的默认轨迹因为默认机制存在丢弃可能。这种场景通常是自己实现一个SendMessageHook把消息关键信息同步写到本地或者外部存储或者由消费者在业务处理成功后主动再发一条“处理完成”的标记消息。本质上是用额外的存储成本换可靠性。5.3 追问三事务消息、顺序消息和延迟消息的轨迹有什么特殊之处这三个场景面试里经常连环追问。先说事务消息半消息的发送阶段会正常触发发送钩子所以你能看到发送轨迹但事务回查、commit、rollback这些内部流程不会额外生成轨迹如果你在轨迹里看到一条事务消息最终没有被消费需要结合业务侧的事务状态去判断。第二个是顺序消息顺序消息消费失败时的默认行为是暂停该队列的后续消费并重试轨迹里会出现连续的SubBefore和SubAfter记录透出的信息是该队列在某个时间段内被阻塞住了。第三个是延迟消息它的发送轨迹显示发送时间和存储时间但如果消息设定了延迟级别SubBefore的时间会明显晚于发送时间这中间的时间差是正常的固有延迟。5.4 追问四轨迹Topic会不会无限膨胀怎么控制轨迹消息也是普通消息RocketMQ对消息的清理策略不是按Topic的TTL来做的而是基于CommitLog文件的保留时间或者磁盘使用率。默认情况下超过保留时间的CommitLog文件会被整体删除里面的轨迹消息自然就没了。所以你其实没法简单地对轨迹Topic单独设置“保留7天、业务Topic保留3天”而是整个Broker有一个统一的文件保留策略。那生产上怎么控制轨迹Topic的膨胀我的做法是三条第一条用自定义轨迹Topic把轨迹数据和业务消息放在不同的Broker组或者不同磁盘上第二条在承接轨迹Topic的Broker上调低文件保留时间第三条如果需要的保留周期比Broker上消息保留周期更长就额外写一个定时任务消费轨迹Topic把轨迹数据转存到外部存储然后允许Broker上的原始轨迹被清理。5.5 面试加分表达给轨迹定一个“身份”如果你能在这个问题里给出一个有总结感的个人观点会明显加分。我的观点是把消息轨迹理解成消息中间件给开发者的一份“体检报告”它在设计上就默认了“可以漏掉几条体检数据但绝不能因为做体检把病人搞死”。这种“诊断数据与业务数据分离”的思路在所有大型系统里几乎都能看到。能用一句话说清楚这个取舍比背一百行源码有效得多。如果你正在准备面试建议不要满足于看这篇文字。本地搭一个单机RocketMQ环境开轨迹发几条正常和异常消息再对着控制台看一遍时间线整个过程半小时左右但你对这个问题的理解会从“知道”变成“见过”。面试时你能聊出那些轨迹状态在真实场景下长什么样这个状态就是最好的答案。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进