- 后端
- 消息队列
- 微服务
- 消息路由
【免费下载链接】CAP
Distributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern
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,它的工作流程分为三步:
- 调用
builder.AddSource(DiagnosticListener.SourceName),把 CAP 的 ActivitySource 名称(DotNetCore.CAP.OpenTelemetry,见 DiagnosticListener.cs)注册进 TracerProvider,使 OpenTelemetry 能够采集该 Source 产生的所有 Activity; - 创建
CapInstrumentation实例,其内部通过DiagnosticSourceSubscriber订阅 CAP 的诊断源(CapInstrumentation.cs); - 通过
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/ErrorPublishMessageStore | Event Persistence: {operation} | Internal |
| 消息发送 | Before/After/ErrorPublish | CAP/{operation}/Publisher | Producer |
| 消息消费存储 | Before/After/ErrorConsume | CAP/{operation}/Subscriber | Consumer |
| 订阅者调用 | Before/After/ErrorSubscriberInvoke | Subscriber 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
相关推荐
Laya 预测链路追踪实战:基于 run_id 与 Hook 的 Span 关联、错误追踪与 OpenTelemetry 集成
Laya 预测链路追踪实战:基于 run_id 与 Hook 的 Span 关联、错误追踪与 OpenTelemetry 集成 导读 Laya 是一个非自回归的
人工智能NLP强化学习minikube 遥测(Telemetry)实战:基于 OpenTelemetry 追踪 `minikube start` 全链路
minikube 遥测(Telemetry)实战:基于 OpenTelemetry 追踪 minikube start 全链路 minikube 内置了基于 O
云原生容器编排CLI开发工具微服务链路追踪实战:基于nerdctl部署Jaeger与OpenTelemetry全链路监控
微服务链路追踪实战:基于nerdctl部署Jaeger与OpenTelemetry全链路监控 在微服务架构中,全链路监控是排查分布式系统问题的关键。但传统部署方
CLI云原生
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考