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

