惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
月光博客
月光博客
MyScale Blog
MyScale Blog
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
爱范儿
爱范儿
P
Proofpoint News Feed
人人都是产品经理
人人都是产品经理
Last Week in AI
Last Week in AI
罗磊的独立博客
G
Google Developers Blog
Y
Y Combinator Blog
博客园 - 【当耐特】
WordPress大学
WordPress大学
大猫的无限游戏
大猫的无限游戏
博客园 - 叶小钗
J
Java Code Geeks
酷 壳 – CoolShell
酷 壳 – CoolShell
V
Visual Studio Blog
美团技术团队
宝玉的分享
宝玉的分享
Jina AI
Jina AI
小众软件
小众软件
T
Tailwind CSS Blog
A
About on SuperTechFans

博客园 - Hey,Coder!

微信小程序 web方案实现调音器功能 systemctl slice配置docker最大占用cpu 内存 pipenv环境变量配置 .net8验证码 数据库差异对比工具 docker部署prometheus grafana pgexporter Alertmanager ubuntu24 换源 ssh配置 ubuntu24 Harbor部署 ubuntu24 kvm部署 cockpit web管理界面 .net api代理 Ubuntu 24.04 DMAR IOMMU 报错修复 uniapp iOS App 上架 vs文件软链接 Docker 网络别名 .net8通用中间件处理公共返回值结构 .net环境下OpenTelemetry Body 读取注意事项 jaeger界面显示异常 .net opentelemetry collector jaeger c# opentelemetry 自定义export c# 生成对象的初始化代码 Elasticsearch ILM 与 Data Stream 概念及实例说明文档 portainer httphelper封装 Ocelot + Consul + SignalR 粘性会话 正交软件架构 docker postgresql17 主从复制 rabbitmq总结与c#示例 docker rabbitmq quartz dashboard docker consul registrator 服务自动注册consul 微服务动态扩容
c# opentelemetry 自定义export
Hey,Coder! · 2026-07-22 · via 博客园 - Hey,Coder!

OpenTelemetry 自定义 Exporter + MQ + 请求体记录

整体架构

.NET App (OTel SDK) → 自定义 Exporter → MQ (RabbitMQ/Kafka) → 消费者 → Elasticsearch
  • 目的:解耦应用与 ES,利用 MQ 削峰填谷,提高稳定性。
  • 优势:应用不直接依赖 ES 可用性,消费者可按需批量写入。

nuget

dotnet add package OpenTelemetry --version 1.17.0
dotnet add package OpenTelemetry.Exporter.OpenTelemetryProtocol --version 1.17.0
dotnet add package OpenTelemetry.Extensions.Hosting --version 1.17.0
dotnet add package OpenTelemetry.Instrumentation.AspNetCore --version 1.17.0
dotnet add package OpenTelemetry.Instrumentation.Http --version 1.17.0

自定义 Exporter 实现

1. 继承 BaseExporter<Activity>

/// <summary>
/// 
/// </summary>
public class RabbitMqTraceExporter : BaseExporter<Activity>
{
    //private readonly IMqPublisher _publisher;  // 你的 MQ 发布接口

    public RabbitMqTraceExporter()
    {
    }

    public override ExportResult Export(in Batch<Activity> batch)
    {
        //关键:抑制自产自销的遥测,避免无限循环[6](@ref)
        using var scope = SuppressInstrumentationScope.Begin();

        try
        {
            foreach (var activity in batch)
            {
                // 将 Activity 转换为可序列化对象
                var spanDoc = new Dictionary<string, object?>
                {
                    ["traceId"] = activity.TraceId.ToHexString(),
                    ["spanId"] = activity.SpanId.ToHexString(),
                    ["parentSpanId"] = activity.ParentSpanId.ToHexString(),
                    ["name"] = activity.DisplayName,
                    ["kind"] = (int)activity.Kind,
                    ["startTimeUnixNano"] = activity.StartTimeUtc.Ticks * 100,
                    ["endTimeUnixNano"] = (activity.StartTimeUtc + activity.Duration).Ticks * 100,
                    ["attributes"] = activity.TagObjects.Select(t => new
                    {
                        key = t.Key,
                        value = new { stringValue = t.Value?.ToString() }
                    }).ToList(),
                    ["resource"] = ParentProvider?.GetResource()?.Attributes
                        .Select(a => new { key = a.Key, value = a.Value }).ToList()
                };

                // 发送到 MQ(注意:SDK 不做重试,需要自己实现)[6](@ref)
                //_publisher.Publish("trace", spanDoc);
            }
            return ExportResult.Success;
        }
        catch (Exception ex)
        {
            // 绝不能抛出异常[6](@ref)
            Console.Error.WriteLine($"MQ export failed: {ex.Message}");
            return ExportResult.Failure;
        }
        return ExportResult.Success;
    }
}

注册函数封装

    public static void AddOT(this WebApplicationBuilder builder, IConfiguration configuration)
    {
        //builder.Services.AddOpenTelemetry()
        //.ConfigureResource(resource => resource.AddService("my-service"))
        //.WithTracing(tracing => tracing
        //    .AddAspNetCoreInstrumentation()
        //    .AddHttpClientInstrumentation()
        //    .AddOtlpExporter(opt =>  // ← 现在能找到了
        //    {
        //        opt.Endpoint = new Uri("https://your-es:9200/_otlp");
        //        opt.Protocol = OtlpExportProtocol.HttpProtobuf;
        //        opt.Headers = "Authorization=ApiKey your_api_key";
        //    }));



        builder.Services.AddOpenTelemetry()
            .ConfigureResource(resource => resource.AddService("my-service"))
            .WithTracing(tracing => tracing
                .AddAspNetCoreInstrumentation(options =>
                {
                    // 1) 过滤:排除健康检查
                    options.Filter = ctx =>
                        !ctx.Request.Path.StartsWithSegments("/health");

                    // 2) 记录异常详情
                    options.RecordException = true;

                    // 3) 入参:EnrichWithHttpRequest
                    options.EnrichWithHttpRequest = async (activity, request) =>
                    {
                        var ctx = request.HttpContext;

                        // —— 基础信息 ——
                        activity.SetTag("request.protocol", request.Protocol);
                        activity.SetTag("request.scheme", request.Scheme);
                        activity.SetTag("request.host", request.Host.ToString());
                        activity.SetTag("request.path", request.Path.ToString());
                        activity.SetTag("request.query_string", request.QueryString.ToString());

                        // —— 客户端 IP ——
                        activity.SetTag("request.client_ip",
                            ctx.Connection.RemoteIpAddress?.ToString());

                        // —— Query 参数(扁平化)——
                        foreach (var q in request.Query)
                        {
                            activity.SetTag($"request.query.{q.Key}", q.Value.ToString());
                        }

                        // —— 指定 Header(注意脱敏)——
                        if (request.Headers.TryGetValue("X-Correlation-Id", out var corr))
                            activity.SetTag("request.correlation_id", corr.ToString());
                        if (request.Headers.TryGetValue("User-Agent", out var ua))
                            activity.SetTag("request.user_agent", ua.ToString());

                        // —— 请求体(仅当存在且体积可控)——
                        if (request.ContentLength > 0 && request.ContentLength <= 1024 * 1024)
                        {
                            try
                            {
                                // 确保已启用缓冲(以防外部未设置)
                                request.EnableBuffering();

                                // 1. 将 Body 内容复制到 MemoryStream
                                 var memoryStream = new MemoryStream();
                                await request.Body.CopyToAsync(memoryStream);
                                memoryStream.Position = 0;

                                // 2. 读取 Body 内容
                                using var reader = new StreamReader(memoryStream, Encoding.UTF8, leaveOpen: true);
                                var body = reader.ReadToEnd();
                                activity.SetTag("request.body", body);

                                // 3. 将 request.Body 替换为 MemoryStream(可 Seek)
                                memoryStream.Position = 0;
                                request.Body = memoryStream; // 替换后,MVC 模型绑定就能正常读取了
                            }
                            catch (Exception ex)
                            {
                                // 兜底,防止读取 Body 导致主程序崩溃
                                activity.SetTag("request.body.error", ex.Message);
                            }
                        }
                    };

                    // 4) 出参:EnrichWithHttpResponse
                    options.EnrichWithHttpResponse = async (activity, response) =>
                    {
                        activity.SetTag("response.status_code", response.StatusCode);
                        activity.SetTag("response.content_type", response.ContentType);

                        if (response.HttpContext.Items.TryGetValue("response_body", out var bodyObj))
                        {
                            var body = bodyObj as string;
                            if (!string.IsNullOrEmpty(body))
                            {
                                // 脱敏
                                //body = MaskSensitiveFields(body);
                                activity.SetTag("response.body", body);
                            }
                        }
                    };

                    // 5) 异常:EnrichWithException
                    options.EnrichWithException = (activity, exception) =>
                    {
                        activity.SetTag("exception.type", exception.GetType().FullName);
                        activity.SetTag("exception.message", exception.Message);
                        // 堆栈通过 RecordException=true 自动作为 ActivityEvent 附加
                    };
                })
                .AddHttpClientInstrumentation()
                .AddHttpClientInstrumentation()
                // 使用工厂函数注入自定义 Exporter,并用 BatchActivityExportProcessor 包装
                .AddProcessor(sp =>
                {
                    //var mqPublisher = sp.GetRequiredService<IMqPublisher>();
                    var exporter = new RabbitMqTraceExporter();

                    // 生产环境必须用 Batch,不要用 Simple[6](@ref)
                    return new BatchActivityExportProcessor(
                        exporter,
                        maxQueueSize: 2048,
                        scheduledDelayMilliseconds: 5000,
                        exporterTimeoutMilliseconds: 30000,
                        maxExportBatchSize: 512);
                }));
    }

注册

builder.AddOT(builder.Configuration);

必须的管道配置(Program.cs)

app.Use((context, next) =>
{
    context.Request.EnableBuffering();
    return next();
});
// 然后才是 UseRouting、UseMiddleware 等

//在routing后面增加

app.UseMiddleware<ResponseCaptureMiddleware>();  // 放在路由之后、端点之前

ResponseCaptureMiddleware

 public class ResponseCaptureMiddleware
 {
     private readonly RequestDelegate _next;

     public ResponseCaptureMiddleware(RequestDelegate next) => _next = next;

     public async Task InvokeAsync(HttpContext context)
     {
         var originalBodyStream = context.Response.Body;

         using var memoryStream = new MemoryStream();
         context.Response.Body = memoryStream;

         await _next(context);

         // 读取响应体
         memoryStream.Position = 0;
         var responseBody = await new StreamReader(memoryStream).ReadToEndAsync();

         // 把响应体挂到 HttpContext.Items,供 Enrich 回调读取
         if (responseBody.Length <= 1024 * 1024)
             context.Items["response_body"] = responseBody;

         // 把流写回原始 Response
         memoryStream.Position = 0;
         await memoryStream.CopyToAsync(originalBodyStream);
         context.Response.Body = originalBodyStream;
     }
 }

性能与安全提醒

  1. 限制 Body 大小ContentLength <= 1MB,防止大文件撑爆 MQ/ES。
  2. 脱敏:在记录前替换敏感字段(密码、Token 等)。
  3. 采样:高流量接口启用采样(如 TraceIdRatioBasedSampler)。
  4. 抑制自监控SuppressInstrumentationScope.Begin() 防止 Exporter 自身产生的遥测被再次采集。
  5. 包版本对齐:所有 OpenTelemetry.* 包版本一致(推荐 1.9+)。