十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

SkyWalking 消息队列(Message Queue)消费性能与消费延迟监控指南

SkyWalking 消息队列(Message Queue)消费性能与消费延迟监控指南 SkyWalking 消息队列Message Queue消费性能与消费延迟监控指南【免费下载链接】skywalkingAPM, Application Performance Monitoring System项目地址: https://gitcode.com/gh_mirrors/sky/skywalking导读本文基于 Apache SkyWalking 官方后端部署文档 docs/en/setup/backend/mq.md系统讲解如何在分布式系统中监控消息队列Message Queue的消费流量与消费延迟。文章将从为什么需要 MQ 消费监控出发深入剖析其底层实现原理——基于 SkyWalking Cross Process Propagation Headers Protocol v3 的 Extension Header Item 实现端到端延迟计算并完整给出默认开箱即用的两类指标、服务/端点两级作用域以及如何通过core.oal自定义扩展更多 MQ 指标。读完本文你将掌握 SkyWalking 监控 MQ 消费性能的完整原理与自定义扩展方法。背景为什么消息队列的消费性能监控如此重要在现代分布式系统中消息队列Message Queue扮演着至关重要的角色。通过将请求解耦为异步消息传递MQ 可以显著缩短阻塞式 RPC同步远程调用的链路长度与响应延迟从而最终改善用户体验。然而这种异步化的处理方式也带来了新的观测挑战在同步 RPC 中请求的耗时天然由调用链完整承载而在异步 MQ 场景下消息从生产者发送到消费者接收之间存在传输、排队、消费处理等多段不可见的时间。如果缺乏对消息队列消费流量和消费延迟的度量就无法判断系统瓶颈究竟出现在消息堆积、网络传输还是消费者处理能力上。正是基于这一需求SkyWalking 从 8.9.0 版本开始利用原生追踪探针native tracing agent配合 Cross Process Propagation Headers Protocol v3 的 Extension Header Item提供了对消息队列系统的性能监控能力。核心原理Extension Header Item 与transmission.latencysw8-x 扩展头的作用SkyWalking 的跨进程上下文传播协议 v3又称 sw8 协议包含两类 Header标准 Header Itemsw8用于最基本的上下文传播携带采样标记、Trace ID、父 Segment ID、父 Span ID、父服务、父服务实例、父端点以及目标地址等 8 个字段Extension Header Itemsw8-x用于高级特性为上下游服务中部署的探针之间提供交互能力字段以-分隔且可扩展。关于 Extension Header Item 的完整协议定义参见 SkyWalking Cross Process Propagation Headers Protocol v3。消费延迟如何被计算sw8-x扩展头中的关键字段之一是客户端发送时间戳the timestamp of sending on the client end。该字段专门用于异步 RPC如 MQ场景一旦设置该时间戳消费者端会自动计算发送与接收之间的延迟并将延迟以键transmission.latency标记tag在 Span 上。也就是说transmission.latency的语义为消息从生产者发送到消费者接收所经过的传输耗时其计算公式即transmission.latency timestamp(received) - timestamp(producing)这一点在 UI 模板中也有明确说明见 general-service.json 中 Message Queue Avg Consuming Latency (ms) 的 tips 文本The avg latency of message consuming. Latency timestamp(received) - timestamp(producing)。底层实现路径从源码层面可以确认该监控链路的具体落点Tag 键定义在 SpanTags.java 中定义了常量TRANSMISSION_LATENCY transmission.latency这是 OAP 侧识别该延迟标记的统一键名Span 类型识别在 RPCAnalysisListener.java 的setPublicAttrs方法中根据span.getSpanLayer()判断请求类型MQ层的 Span 会被设置为RequestType.MQ从而与 HTTP、DATABASE、RPC 等请求类型区分开供后续 OAL 脚本按type RequestType.MQ过滤Tag 透传在同一个方法中Span 上携带的所有 Tag包括transmission.latency会通过sourceBuilder.setTag(tag)写入 Source最终进入指标计算。默认指标开箱即用的两类 MQ 监控指标按官方文档说明SkyWalking 默认提供以下两个指标作用域覆盖服务Service级与端点Endpoint级指标名称含义作用域Message Queue Consuming Count消息队列消费次数消费流量Service / EndpointMessage Queue Avg Consuming Latency消息队列平均消费延迟Service / Endpoint服务级Service指标对应 OAL 脚本位于 core.oalservice_mq_consume_count from(Service.*).filter(type RequestType.MQ).count(); service_mq_consume_latency from((str-long)Service.tag[transmission.latency]).filter(type RequestType.MQ).filter(tag[transmission.latency] ! null).longAvg();service_mq_consume_count对服务维度下所有RequestType.MQ类型的请求做计数反映消息消费的流量吞吐service_mq_consume_latency先将字符串类型的transmission.latencyTag 强制转换为long再对带有该 Tag 的 MQ 请求求longAvg得到服务级的平均消费延迟。filter(tag[transmission.latency] ! null)保证只有携带该延迟标记的请求参与计算。端点级Endpoint指标对应 OAL 脚本位于 core.oalendpoint_mq_consume_latency from((str-long)Endpoint.tag[transmission.latency]).filter(type RequestType.MQ).filter(tag[transmission.latency] ! null).longAvg();注意端点级默认只提供了平均消费延迟指标endpoint_mq_consume_latency端点级的消费次数可直接复用通用的endpoint_cpm每分钟调用量等指标。UI 模板中的对应图表这两类指标在 SkyWalking UI 的默认服务概览模板中已预先配置好见 general-service.json名为Message Queue Consuming Count的折线图 Widget绑定表达式service_mq_consume_countlabel 为consuming_count名为Message Queue Avg Consuming Latency (ms)的折线图 Widget绑定表达式service_mq_consume_latency并带有提示文本 The avg latency of message consuming. Latency timestamp(received) - timestamp(producing)。部署后端后在服务Service页面的概览中即可直接看到这两张 MQ 监控图表。深入 OAL完整的 MQ 访问指标族除了上文默认暴露的两个消费侧指标外当前仓库的 core.oal 中还定义了更完整的 MQ 访问指标族它们基于MQAccess/MQEndpointAccess两个 Source并区分MQOperation.Consume消费与MQOperation.Produce生产两种操作mq_service_consume_cpm from(MQAccess.*).filter(operation MQOperation.Consume).cpm(); mq_service_consume_sla from(MQAccess.*).filter(operation MQOperation.Consume).percent(status true); mq_service_consume_latency from(MQAccess.transmissionLatency).filter(operation MQOperation.Consume).longAvg(); mq_service_consume_percentile from(MQAccess.transmissionLatency).filter(operation MQOperation.Consume).percentile2(10); mq_service_produce_cpm from(MQAccess.*).filter(operation MQOperation.Produce).cpm(); mq_service_produce_sla from(MQAccess.*).filter(operation MQOperation.Produce).percent(status true); mq_endpoint_consume_cpm from(MQEndpointAccess.*).filter(operation MQOperation.Consume).cpm(); mq_endpoint_consume_latency from(MQEndpointAccess.transmissionLatency).filter(operation MQOperation.Consume).longAvg(); mq_endpoint_consume_percentile from(MQEndpointAccess.transmissionLatency).filter(operation MQOperation.Consume).percentile2(10); mq_endpoint_consume_sla from(MQEndpointAccess.*).filter(operation MQOperation.Consume).percent(status true); mq_endpoint_produce_cpm from(MQEndpointAccess.*).filter(operation MQOperation.Produce).cpm(); mq_endpoint_produce_sla from(MQEndpointAccess.*).filter(operation MQOperation.Produce).percent(status true);该指标族覆盖的能力包括指标说明cpm每分钟调用量call per minute衡量消费/生产吞吐sla成功率percent(status true)latency平均传输延迟基于MQAccess.transmissionLatencypercentile延迟百分位percentile2(10)可得出 p50/p75/p90/p95/p99 等多档分位值自定义扩展通过 core.oal 添加更多 MQ 指标官方文档明确指出更多的 MQ 指标可以通过core.oal添加。操作步骤找到 OAL 定义文件后端默认的指标定义位于 oap-server/server-starter/src/main/resources/oal/core.oal在该文件中追加自定义的 OAL 表达式语法参考已有 MQ 指标例如增加服务级 MQ 消费延迟的 p99 分位数service_mq_consume_percentile from((str-long)Service.tag[transmission.latency]).filter(type RequestType.MQ).filter(tag[transmission.latency] ! null).percentile2(10);在 general-service.json 或自定义 UI 模板中为新增指标添加对应的 Widgetexpressions字段填入新的 OAL 指标名如service_mq_consume_percentile重启 OAP Server 使 OAL 定义生效。OAL 语法速览from(Source.*)/from(Source.field)指定指标来源与取值字段.filter(condition)按条件过滤MQ 场景常用type RequestType.MQ、operation MQOperation.Consume/MQOperation.Produce、tag[transmission.latency] ! null聚合函数count()、cpm()、longAvg()、percent(status true)、percentile2(10)类型转换(str-long)Service.tag[transmission.latency]将字符串 Tag 转为长整型参与数值计算。关于 OAL 的完整语法体系可进一步阅读 OAL 设计与概念文档。版本与适用前提该功能自SkyWalking 8.9.0起提供相关的变更记录见 changes-8.9.0.mdAddMessage Queue Consuming Countmetric for MQ consuming service and endpoint 与 AddMessage Queue Avg Consuming Latencymetric for MQ consuming service and endpoint本文所引用的 OAL 定义、UI 模板与源码均来自当前仓库对应 10.x 版本线默认指标以当前仓库实际内容为准该能力依赖 SkyWalking 原生探针native tracing agent上报的 MQ 层 Span 及sw8-x扩展头使用前请确认探针版本已支持对应的 MQ 插件如 Kafka、RocketMQ、RabbitMQ、Pulsar 等与 v3 传播协议若你的应用场景涉及非探针方式接入的虚拟 MQ如 Envoy、Istio 等 Service Mesh 场景可参考 virtual-mq 文档 了解相关指标同样基于MQAccessSource。总结SkyWalking 对消息队列的消费性能监控核心路径为探针在生产者端通过sw8-x扩展头写入发送时间戳 → 消费者端收到消息后计算transmission.latency并标记到 Span → OAP 侧依据 Span 层MQ与 Tag 值通过 core.oal 中的 OAL 表达式聚合出服务级与端点级的消费次数、平均消费延迟指标 → 最终在 UI 的 general-service.json 模板中呈现为可视化图表。默认两个指标开箱即用更丰富的消费/生产吞吐、成功率与延迟分位指标已在MQAccess指标族中定义且用户可随时通过扩展core.oal自定义属于自己的 MQ 监控指标。【免费下载链接】skywalkingAPM, Application Performance Monitoring System项目地址: https://gitcode.com/gh_mirrors/sky/skywalking创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表