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

ASP.NET Core 8 Web API高负载下流式上传不稳定问题求助

高并发流式文件转发时出现Unexpected end of request content的问题排查与解决

问题背景

ASP.NET Core 8 Web API实现流式文件上传,将文件转发至3个服务,采用以下流程:

  • 单线程读取Request.Body分块(80k)后放入3个队列
  • 自定义流通过队列响应Read()方法,无数据时用Monitor.Wait()等待
  • 用HttpClient将自定义流作为请求体发送至目标服务

低并发(3-5个请求)正常,50个并行请求时,接收转发请求的服务端抛出:

Microsoft.AspNetCore.Server.Kestrel.Core.BadHttpRequestException: Unexpected end of request content.

服务运行于OpenShift Docker容器,CPU负载仅2-3核(容器限制16核),无资源耗尽迹象。

可能的原因

  1. 手动同步机制不可靠:Monitor.Wait/Pulse在高并发下易出现信号丢失、虚假唤醒,导致自定义流错误返回0(标记流结束),接收端收到不完整内容。
  2. 同步流与异步模型冲突:自定义流仅实现同步Read方法,而ASP.NET Core、HttpClient优先使用异步流操作,混用同步等待易引发线程池阻塞、流状态异常,甚至请求被提前终止。
  3. 取消令牌未传递:未绑定请求的HttpContext.RequestAborted令牌,当请求因超时/客户端中断被取消时,生产者线程和自定义流未及时终止,导致接收端等待超时或收到不完整数据。
  4. 队列边界处理失误:输入流读取完成后,未正确给所有自定义流发送结束信号,部分流可能无限等待,或错误提前返回0。

解决办法

1. 替换手动队列+Monitor为Channel

.NET Core 3.0+引入的Channel<T>专为异步生产者消费者场景设计,自带线程安全、异步等待机制,彻底避免Monitor的同步问题。

生产者代码(读取Request.Body写入Channel)

// 为每个转发目标创建有界Channel(限制队列长度避免内存溢出)
var channels = Enumerable.Range(0, 3)
    .Select(_ => Channel.CreateBounded<byte[]>(new BoundedChannelOptions(10)
    {
        FullMode = BoundedChannelFullMode.Wait // 队列满时等待,避免丢数据
    }))
    .ToList();

var requestStream = HttpContext.Request.Body;
var buffer = new byte[81920]; // 80k块大小
int bytesRead;
var cancellationToken = HttpContext.RequestAborted;

try
{
    while ((bytesRead = await requestStream.ReadAsync(buffer, cancellationToken)) > 0)
    {
        // 复制当前块(避免buffer被后续读取覆盖)
        var chunk = buffer.AsSpan(0, bytesRead).ToArray();
        // 并行写入所有Channel提升效率
        await Task.WhenAll(channels.Select(c => c.Writer.WriteAsync(chunk, cancellationToken)));
    }
}
finally
{
    // 无论成功还是取消,都标记所有Channel完成
    foreach (var channel in channels)
    {
        channel.Writer.TryComplete(cancellationToken.IsCancellationRequested ? cancellationToken : null);
    }
}

基于Channel的自定义流实现

public class ChannelBackedStream : Stream
{
    private readonly ChannelReader<byte[]> _reader;
    private byte[] _currentChunk;
    private int _currentOffset;
    private readonly CancellationToken _linkedToken;

    public ChannelBackedStream(ChannelReader<byte[]> reader, CancellationToken cancellationToken)
    {
        _reader = reader;
        _linkedToken = cancellationToken;
    }

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

    public override async Task<int> ReadAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken)
    {
        // 合并请求取消令牌和流自身令牌
        using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, _linkedToken);
        var combinedToken = cts.Token;

        while (true)
        {
            combinedToken.ThrowIfCancellationRequested();

            // 处理当前剩余的块数据
            if (_currentChunk != null)
            {
                var copySize = Math.Min(count, _currentChunk.Length - _currentOffset);
                Buffer.BlockCopy(_currentChunk, _currentOffset, buffer, offset, copySize);
                _currentOffset += copySize;

                if (_currentOffset == _currentChunk.Length)
                {
                    _currentChunk = null;
                    _currentOffset = 0;
                }
                return copySize;
            }

            // 等待Channel有数据可读
            if (!await _reader.WaitToReadAsync(combinedToken))
            {
                // Channel已完成,返回0表示流结束
                return 0;
            }

            // 读取下一块数据
            if (_reader.TryRead(out var chunk))
            {
                _currentChunk = chunk;
            }
        }
    }

    // 同步Read委托给异步实现,避免同步异步混用问题
    public override int Read(byte[] buffer, int offset, int count)
    {
        return ReadAsync(buffer, offset, count, CancellationToken.None).GetAwaiter().GetResult();
    }

    // 未实现的方法抛出NotSupportedException
    public override void Flush() => 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) => throw new NotSupportedException();
}

2. 正确传递取消令牌

确保所有异步操作绑定HttpContext.RequestAborted令牌,当请求被取消时,所有相关操作及时终止,避免资源泄漏或无效等待。

3. 优化HttpClient使用

  • 使用IHttpClientFactory创建实例,避免高并发下Socket耗尽
  • 为转发请求设置合理超时:
var httpClient = _httpClientFactory.CreateClient();
httpClient.Timeout = TimeSpan.FromMinutes(5); // 根据业务调整

// 并行发送所有转发请求
var tasks = channels.Select(channel =>
{
    var stream = new ChannelBackedStream(channel.Reader, HttpContext.RequestAborted);
    var request = new HttpRequestMessage(HttpMethod.Post, "https://target-service/api/upload")
    {
        Content = new StreamContent(stream)
        {
            Headers = { ContentType = HttpContext.Request.ContentType } // 保持原Content-Type
        }
    };
    return httpClient.SendAsync(request, HttpContext.RequestAborted);
});

// 根据业务需求选择是否等待所有请求完成
await Task.WhenAll(tasks);

4. 检查Kestrel配置

确保请求体大小、超时配置满足业务需求:

builder.WebHost.ConfigureKestrel(options =>
{
    options.Limits.MaxRequestBodySize = null; // 无限制,或设置合适大小
    options.Limits.KeepAliveTimeout = TimeSpan.FromMinutes(10);
});

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 03:28:11