You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

.NET 9 Azure Function Kafka分布式追踪传播器实现方案咨询

Azure Function + Kafka 分布式追踪实现方案(.NET 9)

问题背景

使用.NET 9开发Azure Function,已完成基础OpenTelemetry配置(追踪数据可写入Seq),需要通过KafkaOutput/KafkaTrigger绑定,在生产者和消费者Function之间延续分布式追踪链路,但Azure绑定未暴露底层Producer/Consumer的Headers操作能力。


实现步骤

一、生产者端(发送Kafka消息的Function)

由于KafkaOutput绑定不支持直接设置消息Headers,我们通过包装业务消息+手动注入追踪上下文的方式传递链路信息:

  1. 定义带追踪上下文的消息包装类
public class KafkaMessageWithTraceContext<T>
{
    // 原始业务数据
    public T Payload { get; set; } = default!;
    // 存储W3C追踪上下文(traceparent、tracestate等)
    public Dictionary<string, string> TraceHeaders { get; set; } = new();
}
  1. 发送消息时注入追踪上下文
    在生成Kafka输出消息前,用OpenTelemetry的TextMapPropagator将当前Activity的上下文注入到包装类的TraceHeaders中:
// 获取默认W3C传播器
var propagator = Propagators.DefaultTextMapPropagator;
var traceHeaders = new Dictionary<string, string>();

// 注入当前追踪上下文到字典
propagator.Inject(
    new PropagationContext(Activity.Current?.Context ?? default, Baggage.Current),
    traceHeaders,
    (dict, key, value) => dict[key] = value);

// 包装业务数据并赋值给Kafka输出绑定
KafkaEventUploadDocument = new[]
{
    new KafkaMessageWithTraceContext<UploadDocumentModel>
    {
        Payload = new UploadDocumentModel 
        { 
            // 填充你的业务数据
        },
        TraceHeaders = traceHeaders
    }
};

二、消费者端(接收Kafka消息的Function)

在消费消息时,从包装类中提取追踪上下文,设置为当前Activity的父上下文,从而延续链路:

  1. 修改KafkaTrigger的消息类型
[FunctionName("ConsumeUploadDocument")]
public async Task Run(
    [KafkaTrigger("%KafkaBrokerList%",
        "%KafkaTopic_UploadDocument%",
        Username = "%KafkaUsername%",
        Password = "%KafkaPassword%",
        AuthenticationMode = BrokerAuthenticationMode.Plain)]
    KafkaMessageWithTraceContext<UploadDocumentModel>[] messages,
    ILogger log)
{
    foreach (var message in messages)
    {
        // 提取追踪上下文
        var propagator = Propagators.DefaultTextMapPropagator;
        var parentContext = propagator.Extract(
            default,
            message.TraceHeaders,
            (dict, key) => dict.TryGetValue(key, out var value) ? value : null);

        // 启动子Activity,关联父追踪上下文
        using var activity = Tracing.Source.StartActivity(
            "Process Upload Document",
            ActivityKind.Consumer,
            parentContext.ActivityContext);

        // 执行业务逻辑
        await ProcessDocument(message.Payload);
    }
}

三、优化OpenTelemetry配置(可选)

添加Confluent Kafka的Instrumentation,自动捕获Kafka操作的追踪事件(需安装NuGet包OpenTelemetry.Instrumentation.ConfluentKafka):

services.AddOpenTelemetry()
    .ConfigureResource(r => r.AddService("MyApplicationName", serviceVersion: "1.0"))
    .WithTracing(tracing =>
    {
        tracing.SetSampler(new AlwaysOnSampler());
        tracing.AddSource("MyApplicationName");
        // 添加Kafka Instrumentation
        tracing.AddConfluentKafkaInstrumentation(o =>
        {
            o.RecordException = true;
        });
        tracing.AddAspNetCoreInstrumentation(o => { o.RecordException = true; });
        tracing.AddHttpClientInstrumentation(o => { o.RecordException = true; });
        tracing.AddEntityFrameworkCoreInstrumentation();
        tracing.AddNpgsql();

        tracing.AddOtlpExporter(name: "OtlpTracing", configure: null);
    });

关键说明

  • 由于Azure Function Kafka绑定封装了底层Producer/Consumer,无法直接操作消息Headers,因此采用业务消息体携带追踪上下文是当前约束下的最优方案。
  • 确保生产者和消费者使用相同的传播器(默认W3C Trace Context),保证上下文解析一致性。
  • 自定义的RootActivity工具类可以正常配合该方案使用,链路会自动关联。

内容的提问来源于stack exchange,提问作者user3079834

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.16 07:58:11