☰
CAP 与 OpenTelemetry 集成实战:基于 AddCapInstrumentation 的消息全链路追踪
2026/10/7 2:05:53 网站建设 项目流程
  • 后端
  • 消息队列
  • 微服务
  • 消息路由

【免费下载链接】CAP

Distributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern

项目地址:https://gitcode.com/gh_mirrors/ca/CAP
点击查看免费下载

CAP 内置了面向 OpenTelemetry 的仪器化(Instrumentation)支持,通过DotNetCore.CAP.OpenTelemetry包即可将消息发布、持久化、消费与订阅者调用等关键环节的跟踪数据自动接入 OpenTelemetry,并配合 Zipkin 等后端完成可视化。本文以当前仓库中的官方中文指南为主体,结合源码实现细节,说明如何安装与配置 CAP Instrumentation、理解其底层数据来源(Diagnostics)与 Span 结构,以及如何利用 Context Propagation 在消息传递过程中保持分布式上下文的连续性,帮助你在微服务场景中快速搭建出可观测的消息链路。

什么是 OpenTelemetry

OpenTelemetry 是工具、API 与 SDK 的集合,用于对软件进行插桩(Instrument),从而生成、收集并导出遥测数据(度量 Metrics、日志 Logs 与跟踪 Traces),帮助你分析软件的性能与行为。它提供了统一的跨语言标准,使不同服务、不同技术栈产生的可观测数据可以被一致地采集和展示。

在 .NET 场景下,OpenTelemetry 提供了成熟的 .NET SDK,你可以在官方入门文档中找到如何在控制台应用或 ASP.NET Core 中使用它。本文将聚焦于如何把 CAP 集成进 OpenTelemetry,而不是重复介绍 OpenTelemetry 本身的基础用法。

CAP 与 OpenTelemetry 的集成原理

CAP 对 OpenTelemetry 的跟踪数据支持不是独立实现的,而是建立在其已有的 诊断(Diagnostics) 机制之上:CAP 在运行的关键节点会通过 .NETDiagnosticSource发出诊断事件,OpenTelemetry 的 CAP Instrumentation 订阅这些事件,并将其转换为标准Activity(即 Span)供 OpenTelemetry 采集。

  • 诊断监听器名称为CapDiagnosticListener,定义于 CapDiagnosticListenerNames.cs;
  • CAP 对外提供的事件包括:消息持久化之前/之后/异常、消息向 MQ 发送之前/之后/异常、消息从 MQ 消费保存之前/之后、订阅者方法执行之前/之后/异常,共 11 类事件名,均以DotNetCore.CAP.前缀定义在同一个文件(CapDiagnosticListenerNames.cs);
  • 这些事件由 IMessageSender.Default.cs 等内部组件在真实的消息发送流程中写入,事件负载(EventData)则定义于 EventData.Cap.P.cs 与 EventData.Cap.S.cs。

因此,只要在 OpenTelemetry 的扩展配置中添加 CAP Instrumentation,它便会自动订阅上述诊断事件并完成跟踪数据的收集,无需在业务代码中手动埋点。

安装 CAP 的 OpenTelemetry 包

首先,在你的项目中安装DotNetCore.CAP.OpenTelemetry包:

dotnet add package DotNetCore.CAP.OpenTelemetry

从项目文件可以看出,该包以net8.0为目标框架,并依赖OpenTelemetry 1.16.0与核心库DotNetCore.CAP(见 DotNetCore.CAP.OpenTelemetry.csproj),其包描述即为 "CAP instrumentation for OpenTelemetry .NET"。

配置 OpenTelemetry 与 CAP Instrumentation

安装完成后,在Startup或Program.cs的依赖注入配置中添加如下代码,即可启用 CAP 的跟踪数据采集:

services.AddOpenTelemetryTracing((builder) => builder .AddAspNetCoreInstrumentation() .AddCapInstrumentation() // <-- 添加这行 .AddZipkinExporter() );

其中:

  • AddAspNetCoreInstrumentation():采集 ASP.NET Core 自身的请求跟踪数据,用于建立从 HTTP 请求到消息事件的根 Span;
  • AddCapInstrumentation():启用 CAP 的消息事件数据采集,即本篇文章的核心扩展方法;
  • AddZipkinExporter():将跟踪数据导出到 Zipkin,你也可以根据实际选型替换为 Jaeger、OTLP 等其他 Exporter。

从源码实现看,AddCapInstrumentation()定义在 TracerProviderBuilder.Extension.cs,它的工作流程分为三步:

  1. 调用builder.AddSource(DiagnosticListener.SourceName),把 CAP 的 ActivitySource 名称(DotNetCore.CAP.OpenTelemetry,见 DiagnosticListener.cs)注册进 TracerProvider,使 OpenTelemetry 能够采集该 Source 产生的所有 Activity;
  2. 创建CapInstrumentation实例,其内部通过DiagnosticSourceSubscriber订阅 CAP 的诊断源(CapInstrumentation.cs);
  3. 通过builder.AddInstrumentation将实例注册到 Provider 生命周期中,应用关闭时会自动调用Dispose清理订阅。

非 ASP.NET Core 环境下的 ActivityListener

AddOpenTelemetryTracing会自动启用必要的监听器;但如果你所在的环境没有这样的框架处理(例如纯控制台应用),则需要手动注册一个ActivityListener来接收 Activity 事件,例如:

ActivitySource.AddActivityListener(new ActivityListener() { ShouldListenTo = _ => true, Sample = (ref ActivityCreationOptions<ActivityContext> _) => ActivitySamplingResult.AllData, ActivityStarted = activity => Console.WriteLine($"{activity.ParentId}:{activity.Id} - Start"), ActivityStopped = activity => Console.WriteLine($"{activity.ParentId}:{activity.Id} - Stop") });

这样即使没有完整的 OpenTelemetry Provider,也可以观察到 CAP 产生的 Activity 生命周期。

CAP 跟踪数据的 Span 结构与语义

订阅到诊断事件后,DiagnosticListener.cs 会把不同阶段的事件映射为不同命名的 Activity,构成一条完整的消息链路:

阶段事件(CapDiagnosticListenerNames)生成的 Span 名称ActivityKind
消息持久化Before/After/ErrorPublishMessageStoreEvent Persistence: {operation}Internal
消息发送Before/After/ErrorPublishCAP/{operation}/PublisherProducer
消息消费存储Before/After/ErrorConsumeCAP/{operation}/SubscriberConsumer
订阅者调用Before/After/ErrorSubscriberInvokeSubscriber Invoke: {methodName}Internal

对应的命名规则定义在 DiagnosticListener.cs:操作名前缀为CAP/,生产者后缀为/Publisher,消费者后缀为/Subscriber。

每个 Span 还携带丰富的语义化标签(Tag),便于在 Zipkin 等后端中检索和定位:

  • 消息发送与消费阶段会标记messaging.system(消息中间件名称)、messaging.message.id、messaging.message.body.size、messaging.destination.name、server.address与server.port等(DiagnosticListener.cs);
  • 消费阶段额外标记messaging.operation.type、messaging.client.id、messaging.consumer.group.name(DiagnosticListener.cs);
  • 订阅者调用阶段标记code.function.name,即被调用的订阅方法名(DiagnosticListener.cs);
  • 每个阶段结束都会写入带耗时的 ActivityEvent(如cap.send.duration、cap.receive.duration、cap.invoke.duration),失败时则通过SetStatus(ActivityStatusCode.Error, ...)与AddException记录异常(DiagnosticListener.cs)。

下图是 CAP 的跟踪数据在 Zipkin 中的一个实际展示:整个 Trace 由根 Span(HTTP 处理)串起消息持久化、CAP/xxx/Publisher生产者发送、CAP/xxx/Subscriber消费者接收以及Subscriber Invoke订阅者调用等多个嵌套 Span,清晰呈现了一次跨进程消息投递的完整耗时分布:

Context Propagation:跨进程传递跟踪上下文

在分布式消息场景中,生产者与消费者往往位于不同进程,要让一条消息的 Trace 在上下游之间保持连续,就必须在发送时把跟踪上下文注入消息头,在接收时再从中恢复。CAP 通过注入traceparent与baggage两个头来完成这一过程:

  • 发送时:在发布消息前,CAP 会把当前 Activity 上下文与 Baggage 通过Propagator.Inject写入消息的 Headers(见 DiagnosticListener.cs);
  • 接收时:在消费与订阅者调用前,通过Propagator.Extract从消息头还原父级 ActivityContext 与 Baggage,并以它作为父上下文创建新的 Span(见 DiagnosticListener.cs 与 DiagnosticListener.cs)。

CAP 使用Propagators.DefaultTextMapPropagator(见 DiagnosticListener.cs),默认情况下它同时包含TraceContextPropagator与BaggagePropagator。如果你希望关闭 Baggage 的传播,可以在客户端程序中显式覆盖默认传播器,例如:

OpenTelemetry.Sdk.SetDefaultTextMapPropagator( new TraceContextPropagator());

这样消息头中将只保留traceparent,不再传递baggage,适用于对传递数据敏感度有要求的场景。

小结与延伸阅读

通过DotNetCore.CAP.OpenTelemetry包与一行AddCapInstrumentation(),CAP 即可把消息生命周期中的持久化、发送、消费与订阅者调用等关键节点自动接入 OpenTelemetry,配合 Zipkin 等后端实现跨进程的消息链路追踪;其数据源头是 CAP 的 Diagnostics 诊断事件,上下文则通过traceparent/baggage头跨进程传播。关于诊断事件与度量指标的更完整说明,可进一步阅读 CAP 诊断(Diagnostics)指南;若需了解各存储与传输组件的配置方式,可参阅 用户指南目录。

  • 后端
  • 消息队列
  • 微服务
  • 消息路由

【免费下载链接】CAP

Distributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern

项目地址:https://gitcode.com/gh_mirrors/ca/CAP
点击查看免费下载

相关推荐

上一篇:Shell命令别名导出:gh_mirrors/sh1/sh中的别名共享机制
下一篇:Jukebox模型剪枝:通道剪枝与权重稀疏化优化实践

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询