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

推荐订阅源

MongoDB | Blog
MongoDB | Blog
Recorded Future
Recorded Future
Jina AI
Jina AI
The Register - Security
The Register - Security
Threat Intelligence Blog | Flashpoint
Threat Intelligence Blog | Flashpoint
月光博客
月光博客
博客园 - 三生石上(FineUI控件)
F
Fortinet All Blogs
人人都是产品经理
人人都是产品经理
S
SegmentFault 最新的问题
Apple Machine Learning Research
Apple Machine Learning Research
L
LangChain Blog
Y
Y Combinator Blog
H
Hackread – Cybersecurity News, Data Breaches, AI and More
freeCodeCamp Programming Tutorials: Python, JavaScript, Git & More
GbyAI
GbyAI
The GitHub Blog
The GitHub Blog
Vercel News
Vercel News
博客园 - 【当耐特】
雷峰网
雷峰网
The Cloudflare Blog
阮一峰的网络日志
阮一峰的网络日志
aimingoo的专栏
aimingoo的专栏
云风的 BLOG
云风的 BLOG
I
InfoQ
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
Google DeepMind News
Google DeepMind News
Security Latest
Security Latest
有赞技术团队
有赞技术团队
L
Lohrmann on Cybersecurity
P
Proofpoint News Feed
cs.CV updates on arXiv.org
cs.CV updates on arXiv.org
The Last Watchdog
The Last Watchdog
P
Privacy & Cybersecurity Law Blog
Scott Helme
Scott Helme
Google Online Security Blog
Google Online Security Blog
WordPress大学
WordPress大学
Hacker News - Newest:
Hacker News - Newest: "LLM"
NISL@THU
NISL@THU
cs.CL updates on arXiv.org
cs.CL updates on arXiv.org
B
Blog RSS Feed
Cyberwarzone
Cyberwarzone
K
Kaspersky official blog
F
Full Disclosure
Martin Fowler
Martin Fowler
Spread Privacy
Spread Privacy
D
Docker
C
Cisco Blogs
www.infosecurity-magazine.com
www.infosecurity-magazine.com
H
Hacker News: Front Page

博客园 - Hey,Coder!

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 微服务动态扩容 ubuntu24.04 安装docker RAG 检索增强生成 软考 - 架构设计师 知识点总结 c# MailKit3.4.3 发送邮件、附件 nginx 反向代理postgresql c# 信号量 elsa 3.5 中间件记录最后执行的节点数据 xp密钥 c# System.Text.Json 反序列化Dictionary<string,object>时未转换基础类型的处理方法 docker配置代理 openclaw 接入 LMStudio的模型服务 wpf canvas 移动 缩放 windows平台openclaw搭建 wpf scrollerview触摸滚动 c# 动态切换sqlsugar连接字符串 c# Elastic.Clients.Elasticsearch 动态查询 c# es 封装Elastic.Clients.Elasticsearch c# scrollerview滚动到指定元素位置 WPF 重写Expander 基于elsa工作流封装一套变量、组件的体系 常用活动重写 DynamicExpresso 轻量级动态表达式库 c# elsa 3.5.2 程序化工作流常用功能及自定义中间件 leaflet docker 连接老版本的sqlserver ssl错误 c# quartz 动态创建任务 控制任务运行和停止 依赖注入 cefsharp 模拟点击 批量重置数据库的数据 docker 复制远程镜像本地并创建容器 c# newtonsoft dynamic
.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 是否已被其他中间件替换

九、参考资源