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

YARP中间件读取HTTPS响应体写入ElasticSearch遇阻求助

问题分析与解决方案

你的核心问题在于错误地替换了整个IHttpResponseBodyFeature,而非包装原始特性,导致响应数据未正确流向客户端,且未处理响应压缩的情况,最终出现乱码。同时,你没有将捕获的数据回写至原始响应流,会导致客户端收不到完整响应。

关键问题点

  1. 直接替换IHttpResponseBodyFeature为自定义的BodyReader,切断了原始响应流的传递,YARP的输出只会写入你的MemoryStream,客户端无法收到响应。
  2. 未处理响应的Content-Encoding头(如gzip/deflate),读取的是压缩后的字节流,自然显示为乱码。
  3. 自定义BodyReader未实现完整的管道写入逻辑,忽略了PipeWriter的核心作用。

正确实现思路

遵循YARP文档的流式包装建议:装饰原始IHttpResponseBodyFeature,在数据流经时捕获副本,同时保持原始流的正常传输。这样既避免全量缓冲的开销,又能获取响应内容。

步骤1:实现响应捕获的包装类

创建一个装饰器类,包装原始的IHttpResponseBodyFeature,通过双工流同时向原始响应流和捕获流写入数据:

internal class CapturingResponseBodyFeature : IHttpResponseBodyFeature
{
    private readonly IHttpResponseBodyFeature _originalFeature;
    private readonly MemoryStream _capturedStream = new MemoryStream();
    private PipeWriter _capturingWriter;

    public CapturingResponseBodyFeature(IHttpResponseBodyFeature originalFeature)
    {
        _originalFeature = originalFeature;
    }

    public Stream Stream => throw new NotSupportedException("请使用Writer而非Stream进行捕获");

    public PipeWriter Writer
    {
        get
        {
            if (_capturingWriter == null)
            {
                var originalWriter = _originalFeature.Writer;
                _capturingWriter = PipeWriter.Create(new DuplexStream(originalWriter, _capturedStream));
            }
            return _capturingWriter;
        }
    }

    public async Task CompleteAsync()
    {
        if (_capturingWriter != null)
        {
            await _capturingWriter.FlushAsync();
        }
        await _originalFeature.CompleteAsync();
    }

    public void DisableBuffering() => _originalFeature.DisableBuffering();

    public Task SendFileAsync(string path, long offset, long? count, CancellationToken cancellationToken = default)
    {
        // 如需捕获SendFile传输的文件内容,需额外处理;否则直接调用原始方法
        return _originalFeature.SendFileAsync(path, offset, count, cancellationToken);
    }

    public Task StartAsync(CancellationToken cancellationToken = default) => _originalFeature.StartAsync(cancellationToken);

    // 获取捕获的响应字节数组
    public byte[] GetCapturedBytes()
    {
        _capturedStream.Position = 0;
        return _capturedStream.ToArray();
    }

    // 双工流:同时写入原始响应流和捕获流
    private class DuplexStream : Stream
    {
        private readonly PipeWriter _originalWriter;
        private readonly Stream _capturedStream;

        public DuplexStream(PipeWriter originalWriter, Stream capturedStream)
        {
            _originalWriter = originalWriter;
            _capturedStream = capturedStream;
        }

        public override bool CanRead => false;
        public override bool CanSeek => false;
        public override bool CanWrite => true;
        public override long Length => throw new NotSupportedException();
        public override long Position { get => throw new NotSupportedException(); set => throw new NotSupportedException(); }

        public override void Flush()
        {
            _capturedStream.Flush();
            _originalWriter.FlushAsync().AsTask().Wait();
        }

        public override int Read(byte[] buffer, int offset, int count) => throw new NotSupportedException();
        public override long Seek(long offset, SeekOrigin origin) => throw new NotSupportedException();
        public override void SetLength(long value) => throw new NotSupportedException();

        public override void Write(byte[] buffer, int offset, int count)
        {
            _capturedStream.Write(buffer, offset, count);
            _originalWriter.Write(buffer, offset, count);
        }

        public override async Task WriteAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken)
        {
            await _capturedStream.WriteAsync(buffer, offset, count, cancellationToken);
            await _originalWriter.WriteAsync(buffer, offset, count, cancellationToken);
        }
    }
}

步骤2:在中间件中使用捕获类

替换原始特性为自定义捕获类,处理响应压缩,最后将数据推送至ElasticSearch:

public class ResponseCaptureMiddleware
{
    private readonly RequestDelegate _next;

    public ResponseCaptureMiddleware(RequestDelegate next)
    {
        _next = next;
    }

    public async Task InvokeAsync(HttpContext context)
    {
        var originalFeature = context.Features.Get<IHttpResponseBodyFeature>();
        var capturingFeature = new CapturingResponseBodyFeature(originalFeature);
        context.Features.Set<IHttpResponseBodyFeature>(capturingFeature);

        try
        {
            await _next(context);

            // 获取捕获的响应字节
            var capturedBytes = capturingFeature.GetCapturedBytes();
            string responseContent = null;

            // 处理响应压缩,避免乱码
            var contentEncoding = context.Response.Headers.ContentEncoding.ToString();
            if (!string.IsNullOrEmpty(contentEncoding))
            {
                using var decompressedStream = new MemoryStream();
                using var compressedStream = new MemoryStream(capturedBytes);
                switch (contentEncoding.ToLowerInvariant())
                {
                    case "gzip":
                        using var gzip = new GZipStream(compressedStream, CompressionMode.Decompress);
                        await gzip.CopyToAsync(decompressedStream);
                        break;
                    case "deflate":
                        using var deflate = new DeflateStream(compressedStream, CompressionMode.Decompress);
                        await deflate.CopyToAsync(decompressedStream);
                        break;
                    default:
                        decompressedStream.Write(capturedBytes, 0, capturedBytes.Length);
                        break;
                }
                decompressedStream.Position = 0;
                using var sr = new StreamReader(decompressedStream, Encoding.UTF8);
                responseContent = await sr.ReadToEndAsync();
            }
            else
            {
                // 无压缩,直接转换为字符串
                responseContent = Encoding.UTF8.GetString(capturedBytes);
            }

            // 推送请求和响应数据至ElasticSearch
            // PushToElasticSearch(context.Request, responseContent);

        }
        finally
        {
            // 恢复原始响应特性
            context.Features.Set(originalFeature);
        }
    }
}

步骤3:注册中间件

在Program.cs中注册中间件(注意要放在YARP中间件之前):

app.UseMiddleware<ResponseCaptureMiddleware>();
app.MapReverseProxy();

补充说明

  • 请求体的捕获可以用同样的思路:包装HttpRequest.Body为双工流,同时读取和传递请求数据。
  • 若响应包含大文件,建议对捕获的字节流做分片处理,避免内存占用过高。
  • 中间件方式做日志/数据推送是合理的,只要采用流式包装而非全量缓冲,就不会引入显著的性能开销。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 21:35:20