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

推荐订阅源

IT之家
IT之家
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
A
About on SuperTechFans
博客园 - 聂微东
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
freeCodeCamp Programming Tutorials: Python, JavaScript, Git & More
B
Blog RSS Feed
U
Unit 42
Stack Overflow Blog
Stack Overflow Blog
Recent Announcements
Recent Announcements
雷峰网
雷峰网
罗磊的独立博客
Microsoft Security Blog
Microsoft Security Blog
Hugging Face - Blog
Hugging Face - Blog
L
LangChain Blog
人人都是产品经理
人人都是产品经理
The GitHub Blog
The GitHub Blog
F
Fortinet All Blogs
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
H
Help Net Security
P
Proofpoint News Feed
The Cloudflare Blog
D
Docker
大猫的无限游戏
大猫的无限游戏

博客园 - 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通用中间件处理公共返回值结构 jaeger界面显示异常 .net opentelemetry collector jaeger c# opentelemetry 自定义export 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 微服务动态扩容
.net环境下OpenTelemetry Body 读取注意事项
Hey,Coder! · 2026-07-25 · via 博客园 - Hey,Coder!

OpenTelemetry Body 采集与注意事项

测试代码

https://gitee.com/hey_hh/opentele-test.git

一、问题背景

在使用 OpenTelemetry 的 AddAspNetCoreInstrumentation 时,通过 Enrich 回调读取请求/响应体是常见需求。但不当的读取方式会导致并发读取冲突,造成 Body 数据截断或丢失。本文档基于实际代码总结 OTel Body 采集的最佳实践。


二、核心问题:async void + 同步委托 = 并发读取

2.1 委托类型定义

OpenTelemetry 的 EnrichWithHttpRequest 签名为:

// 同步委托(Action 类型)
Action<Activity, HttpRequest> EnrichWithHttpRequest { get; set; }
Action<Activity, HttpResponse> EnrichWithHttpResponse { get; set; }

2.2 async lambda 的陷阱

// ❌ 错误:同步委托 + async lambda → 编译为 async void
options.EnrichWithHttpRequest = async (activity, request) =>
{
    var body = await new StreamReader(request.Body).ReadToEndAsync();
    activity.SetTag("request.body", body);
};

执行时序问题

时间线                          行为
────────────────────────────────────────────────
 t0   OTel 框架调用 EnrichWithHttpRequest
 t1   async void 开始执行
 t2   await ReadToEndAsync() → 返回未完成的 Task
 t3   async void 返回 ← OTel 认为回调完毕
 t4   OTel 继续执行中间件管道
 ────────── 两条执行路径并发 ──────────
 t5a  中间件/控制器读取 Body(模型绑定等)
 t5b  async void 后台继续 await 读取 Body
      ⚠️ 两个读者竞争同一个 Stream!

2.3 Stream 不支持并发读取

HttpRequest.Body 底层是 Kestrel 的 HttpRequestStream,使用单一读取游标:

Stream 内部:
  ┌────────────────────────────────────────┐
  │  缓冲区              读取游标          │
  │  [ 数据 ]      ↑      [数据未读]       │
  └────────────────────────────────────────┘

  读者 A: 游标推进
  读者 B: 同时推进 → 数据错乱或异常

三、解决方案:合并中间件统一采集

3.1 设计思路

将请求体读取和响应体拦截合并到一个中间件中,确保:

  • 单一读取方,消除并发冲突
  • 在 OTel Instrumentation 之前完成 Body 读取
  • 通过 HttpContext.Items 传递数据给 Enrich 回调

3.2 完整实现代码

// APMBodyCaptureMiddleware.cs
using Microsoft.AspNetCore.Http;
using System.Text;

namespace Boaway.Log.Common.APM;

public class APMBodyCaptureMiddleware
{
    public const string RequestBodyKey = "request_body";
    public const string RequestBodyBytesKey = "request_body_bytes";
    public const string ResponseBodyKey = "response_body";

    private static readonly string[] MethodsWithBody = { "POST", "PUT", "PATCH", "DELETE" };
    private readonly RequestDelegate _next;

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

    public async Task InvokeAsync(HttpContext context)
    {
        var request = context.Request;
        var captureRequestBody = MethodsWithBody.Contains(request.Method, StringComparer.OrdinalIgnoreCase)
            && request.Body != null
            && request.Body.CanRead;

        Stream? requestMemStream = null;
        Stream? originalRequestStream = null;
        Stream? requestBodyStream = null;

        // ── 阶段 1:读取请求体 ──
        if (captureRequestBody)
        {
            request.EnableBuffering();
            originalRequestStream = request.Body;
            requestBodyStream = request.Body;
            requestMemStream = new MemoryStream();

            try
            {
                await requestBodyStream!.CopyToAsync(requestMemStream);
                requestMemStream.Position = 0;

                using var reader = new StreamReader(requestMemStream, Encoding.UTF8, leaveOpen: true);
                var body = await reader.ReadToEndAsync();

                if (body.Length > 0 && body.Length <= 1024 * 1024)
                {
                    context.Items[RequestBodyKey] = body;
                    context.Items[RequestBodyBytesKey] = requestMemStream.Length;
                }

                requestMemStream.Position = 0;
                request.Body = requestMemStream;
            }
            catch
            {
                requestMemStream.Dispose();
                throw;
            }
        }

        // ── 阶段 2:拦截响应体 ──
        var originalResponseStream = context.Response.Body;
        using var responseMemStream = new MemoryStream();
        context.Response.Body = responseMemStream;

        try
        {
            await _next(context);

            responseMemStream.Position = 0;
            var responseBody = await new StreamReader(responseMemStream, leaveOpen: true).ReadToEndAsync();

            if (responseBody.Length > 0 && responseBody.Length <= 1024 * 1024)
                context.Items[ResponseBodyKey] = responseBody;

            responseMemStream.Position = 0;
            await responseMemStream.CopyToAsync(originalResponseStream);
        }
        finally
        {
            context.Response.Body = originalResponseStream;

            if (captureRequestBody && requestMemStream != null)
            {
                request.Body = originalRequestStream!;
                requestMemStream.Dispose();
            }
        }
    }
}

3.3 执行流程图

┌─ 请求体阶段 ──────────────────────────────────────────┐
│ 1. 判断是否需要采集(POST/PUT/PATCH/DELETE + Body 可读)│
│ 2. request.EnableBuffering()                           │
│ 3. 读取 request.Body → MemoryStream                   │
│ 4. 存入 context.Items[RequestBodyKey]                 │
│ 5. 替换 request.Body 为 MemoryStream(控制器可重读)   │
│ 6. 替换 response.Body 为 MemoryStream(拦截响应)     │
└─────────────────────────────────────────────────────────┘
                          ↓
                    await _next(context)
                          ↓
┌─ 响应体阶段 ──────────────────────────────────────────┐
│ 7. 读取 response MemoryStream → 存入 Items            │
│ 8. 写回原始 response.Body(发送给客户端)              │
│ 9. finally: 恢复原始 request.Body + 释放资源           │
└─────────────────────────────────────────────────────────┘

3.4 关键设计点

请求体读取:独占访问 + 可重读替换

// EnableBuffering 确保流支持 Position 重置
request.EnableBuffering();

// 读取到内存,独占访问
await requestBodyStream.CopyToAsync(requestMemStream);

// 替换为内存流,控制器可正常读取
request.Body = requestMemStream;

响应体拦截:流替换 + 双重 Position 重置

// 替换响应流,拦截控制器输出
context.Response.Body = responseMemStream;

// 第一次重置:读取响应体
responseMemStream.Position = 0;
var responseBody = await reader.ReadToEndAsync();

// 第二次重置:写回给客户端
responseMemStream.Position = 0;
await responseMemStream.CopyToAsync(originalResponseStream);

异常安全:finally 保证资源恢复

try
{
    await _next(context);
    // ... 读取和写回逻辑
}
finally
{
    // 无论成功或异常,都恢复原始流
    context.Response.Body = originalResponseStream;
    if (captureRequestBody && requestMemStream != null)
    {
        request.Body = originalRequestStream!;
        requestMemStream.Dispose();
    }
}

大小限制:防止内存溢出

if (body.Length > 0 && body.Length <= 1024 * 1024) // 1MB
{
    context.Items[RequestBodyKey] = body;
}

四、在 OTel Enrich 回调中读取

4.1 同步读取 Items(不触碰流)

// APMStartup.cs
public static void AddOT(this WebApplicationBuilder builder)
{
    builder.Services.AddOpenTelemetry()
    .WithTracing(tracing => tracing
        .AddAspNetCoreInstrumentation(options =>
        {
            options.EnrichWithHttpRequest = (activity, request) =>
            {
                // 请求元数据(不涉及 Body 流)
                activity.SetTag("request.protocol", request.Protocol);
                activity.SetTag("request.path", request.Path.ToString());
                // ... 其他元数据
            };

            options.EnrichWithHttpResponse = (activity, response) =>
            {
                activity.SetTag("response.status_code", response.StatusCode);
                activity.SetTag("response.content_type", response.ContentType);

                var ctx = response.HttpContext;

                // 从 Items 读取请求体(同步,无竞态)
                if (ctx.Items.TryGetValue(APMBodyCaptureMiddleware.RequestBodyKey, out var bodyObj)
                    && bodyObj is string body && !string.IsNullOrEmpty(body))
                {
                    activity.SetTag("request.body", body);

                    if (ctx.Items.TryGetValue(APMBodyCaptureMiddleware.RequestBodyBytesKey, out var bytesObj))
                        activity.SetTag("request.body.bytes_read", Convert.ToInt64(bytesObj));

                    activity.SetTag("request.body.string_length", body.Length);
                }

                // 从 Items 读取响应体
                if (ctx.Items.TryGetValue(APMBodyCaptureMiddleware.ResponseBodyKey, out var respBodyObj)
                    && respBodyObj is string respBody && !string.IsNullOrEmpty(respBody))
                {
                    activity.SetTag("response.body", respBody);
                }
            };
        })
        .AddConsoleExporter()
    );
}

4.2 注册顺序

// Program.cs
var app = builder.Build();

// 1. 最先注册 Body 采集中间件(在 OTel Instrumentation 之前)
app.UseMiddleware<APMBodyCaptureMiddleware>();

// 2. 其他中间件
app.UseSwagger();
app.UseAuthorization();

// 3. 控制器
app.MapControllers();
app.Run();

五、async 使用规则

5.1 禁止在同步委托中使用 async lambda

// ❌ Action<T> + async lambda = async void,调用方无法 await
options.EnrichWithHttpRequest = async (activity, request) => { ... };

// ✅ 同步 lambda,从 Items 读取
options.EnrichWithHttpRequest = (activity, request) =>
{
    var body = request.HttpContext.Items[APMBodyCaptureMiddleware.RequestBodyKey] as string;
    activity.SetTag("request.body", body);
};

5.2 委托类型判断

委托签名 async lambda 编译结果 调用方可等待
Action<T> async void
Func<T, Task> async Task

规则:只要委托是 Action 类型,就不能使用 async lambda。

5.3 如果必须读取流

// ⚠️ 次选方案:同步读取已缓冲的流
options.EnrichWithHttpRequest = (activity, request) =>
{
    if (request.Body.CanSeek) request.Body.Position = 0;
    using var reader = new StreamReader(request.Body);
    var body = reader.ReadToEnd(); // 同步读取
    activity.SetTag("request.body", body);
};

仅适用于流已被 EnableBuffering() 缓冲且有 CanSeek 的情况。


六、防御性编程

6.1 检查清单

  • EnrichWithHttpRequest / EnrichWithHttpResponse 是否为同步 lambda?
  • APMBodyCaptureMiddleware 是否在 OTel Instrumentation 之前注册?
  • Body 读取后是否存入 HttpContext.Items
  • 异常路径下是否通过 finally 恢复原始流?
  • Body 大小是否限制在 1MB 以内?

6.2 Content-Type 过滤

// 只记录文本类 Body,跳过二进制
var contentType = request.ContentType;
if (contentType != null &&
    (contentType.Contains("json") ||
     contentType.Contains("xml") ||
     contentType.StartsWith("text/")))
{
    // 采集 Body
}

6.3 敏感数据脱敏

if (body.Contains("password") || body.Contains("token"))
    body = RedactSensitiveData(body);

七、测试策略

7.1 关键测试场景

场景 描述 目的
Content-Length 小请求 < 1KB JSON body 基础验证
Content-Length 大请求 ~2KB JSON body 缓冲机制
Chunked 小分块 100B chunk 流式读取
Chunked 大分块 2KB chunk,~4KB payload 大数据量
Chunked 带延时 50B chunk + 100ms delay 慢速网络竞态放大
并发请求 5 个并行 POST 并发安全
连续请求 20 次顺序 POST 稳定性
超大 payload 100KB+ Chunked 内存管理

7.2 ChunkedContent 实现

class ChunkedContent : HttpContent
{
    private readonly byte[] _data;
    private readonly int _chunkSize;
    private readonly int _delayMs;

    public ChunkedContent(string content, int chunkSize = 1024, int delayMs = 0)
    {
        _data = Encoding.UTF8.GetBytes(content);
        _chunkSize = chunkSize;
        _delayMs = delayMs;
        Headers.ContentType = new MediaTypeHeaderValue("application/json");
    }

    protected override async Task SerializeToStreamAsync(Stream stream, TransportContext context)
    {
        int offset = 0;
        while (offset < _data.Length)
        {
            int chunkLen = Math.Min(_chunkSize, _data.Length - offset);
            await stream.WriteAsync(_data, offset, chunkLen);
            offset += chunkLen;
            if (_delayMs > 0) await Task.Delay(_delayMs);
            await stream.FlushAsync();
        }
    }

    protected override bool TryComputeLength(out long length)
    {
        length = -1;
        return false;
    }
}

7.3 验证点

Activity:
  ├── request.body = {"date":"2024-01-01","temperatureC":25,"summary":"Warm"}
  ├── request.body.string_length = 72
  ├── request.body.bytes_read = 72
  ├── response.status_code = 200
  ├── response.body = {"date":"2024-01-01","temperatureC":25,"summary":"Warm"}
  └── response.content_type = application/json; charset=utf-8

确保 request.bodyresponse.body 内容完整,string_length 与实际长度一致。


八、常见错误排查

8.1 InvalidOperationException: Concurrent reads or writes are not supported.

原因:两个代码路径同时读取 request.Body

排查

  1. 检查是否在 Enrich 回调中使用了 async lambda
  2. 检查是否有其他中间件/过滤器也读取了 Body
  3. 确认 APMBodyCaptureMiddleware 在 OTel 之前注册

8.2 Body 内容为空或不完整

原因:Body 被提前消费,后续读取时流已到达末尾。

排查

  1. 检查是否在读取前调用了 request.EnableBuffering()
  2. 检查读取后是否重置了 Position = 0
  3. 检查 request.Body 是否被正确替换

8.3 ObjectDisposedException: Cannot access a disposed object.

原因:Body 读取后流被释放,后续代码尝试访问。

排查

  1. 检查 finally 块是否正确恢复了原始流
  2. 检查内存流是否使用 leaveOpen: true

8.4 客户端收不到响应

原因:响应体写回失败。

排查

  1. 检查是否在读取后重置了 memoryStream.Position = 0
  2. 检查 CopyToAsync 是否有异常
  3. 检查 Response.Body 是否已被其他中间件替换

九、参考资源