




















https://gitee.com/hey_hh/opentele-test.git
在使用 OpenTelemetry 的 AddAspNetCoreInstrumentation 时,通过 Enrich 回调读取请求/响应体是常见需求。但不当的读取方式会导致并发读取冲突,造成 Body 数据截断或丢失。本文档基于实际代码总结 OTel Body 采集的最佳实践。
OpenTelemetry 的 EnrichWithHttpRequest 签名为:
// 同步委托(Action 类型)
Action<Activity, HttpRequest> EnrichWithHttpRequest { get; set; }
Action<Activity, HttpResponse> EnrichWithHttpResponse { get; set; }
// ❌ 错误:同步委托 + 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!
HttpRequest.Body 底层是 Kestrel 的 HttpRequestStream,使用单一读取游标:
Stream 内部:
┌────────────────────────────────────────┐
│ 缓冲区 读取游标 │
│ [ 数据 ] ↑ [数据未读] │
└────────────────────────────────────────┘
读者 A: 游标推进
读者 B: 同时推进 → 数据错乱或异常
将请求体读取和响应体拦截合并到一个中间件中,确保:
HttpContext.Items 传递数据给 Enrich 回调// 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();
}
}
}
}
┌─ 请求体阶段 ──────────────────────────────────────────┐
│ 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 + 释放资源 │
└─────────────────────────────────────────────────────────┘
// EnableBuffering 确保流支持 Position 重置
request.EnableBuffering();
// 读取到内存,独占访问
await requestBodyStream.CopyToAsync(requestMemStream);
// 替换为内存流,控制器可正常读取
request.Body = requestMemStream;
// 替换响应流,拦截控制器输出
context.Response.Body = responseMemStream;
// 第一次重置:读取响应体
responseMemStream.Position = 0;
var responseBody = await reader.ReadToEndAsync();
// 第二次重置:写回给客户端
responseMemStream.Position = 0;
await responseMemStream.CopyToAsync(originalResponseStream);
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;
}
// 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()
);
}
// Program.cs
var app = builder.Build();
// 1. 最先注册 Body 采集中间件(在 OTel Instrumentation 之前)
app.UseMiddleware<APMBodyCaptureMiddleware>();
// 2. 其他中间件
app.UseSwagger();
app.UseAuthorization();
// 3. 控制器
app.MapControllers();
app.Run();
// ❌ 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);
};
| 委托签名 | async lambda 编译结果 | 调用方可等待 |
|---|---|---|
Action<T> |
async void |
❌ |
Func<T, Task> |
async Task |
✅ |
规则:只要委托是 Action 类型,就不能使用 async lambda。
// ⚠️ 次选方案:同步读取已缓冲的流
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 的情况。
EnrichWithHttpRequest / EnrichWithHttpResponse 是否为同步 lambda?APMBodyCaptureMiddleware 是否在 OTel Instrumentation 之前注册?HttpContext.Items?finally 恢复原始流?// 只记录文本类 Body,跳过二进制
var contentType = request.ContentType;
if (contentType != null &&
(contentType.Contains("json") ||
contentType.Contains("xml") ||
contentType.StartsWith("text/")))
{
// 采集 Body
}
if (body.Contains("password") || body.Contains("token"))
body = RedactSensitiveData(body);
| 场景 | 描述 | 目的 |
|---|---|---|
| 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 | 内存管理 |
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;
}
}
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.body 和 response.body 内容完整,string_length 与实际长度一致。
InvalidOperationException: Concurrent reads or writes are not supported.原因:两个代码路径同时读取 request.Body。
排查:
async lambdaAPMBodyCaptureMiddleware 在 OTel 之前注册原因:Body 被提前消费,后续读取时流已到达末尾。
排查:
request.EnableBuffering()Position = 0request.Body 是否被正确替换ObjectDisposedException: Cannot access a disposed object.原因:Body 读取后流被释放,后续代码尝试访问。
排查:
finally 块是否正确恢复了原始流leaveOpen: true原因:响应体写回失败。
排查:
memoryStream.Position = 0CopyToAsync 是否有异常Response.Body 是否已被其他中间件替换此内容由惯性聚合(RSS阅读器)自动聚合整理,仅供阅读参考。 原文来自 — 版权归原作者所有。