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

.NET中间件实现请求合并(Request Coalescing)方案咨询

.NET 请求合并(Request Coalescing)中间件实现方案

问题分析

你遇到的同路径请求串行触发中间件的问题,大概率是ASP.NET Core的Session串行处理机制导致的——默认情况下,同一Session的请求会被串行执行(因为Session存储是线程安全的,同一用户的请求会等待前一个请求释放Session锁)。如果你的测试是用同一个浏览器(携带相同SessionId)发起请求,就会出现这种现象;换不同浏览器或清除Cookie后,请求就会并行触发中间件。

接下来解决这个问题,并实现真正的请求合并逻辑:

步骤1:解决串行请求问题

如果你的项目不需要Session,直接移除AddSession()和UseSession()的注册即可。如果需要保留Session,可在标记了[Squashable]的端点上添加[SessionState(SessionStateBehavior.Disabled)]特性(需引用Microsoft.AspNetCore.Session包),禁用该端点的Session锁定:

using Microsoft.AspNetCore.Session;

[HttpGet]
[Squashable]
[SessionState(SessionStateBehavior.Disabled)]
public async Task<IActionResult> GetAsync()
{
    await Task.Delay(10000);
    return Ok();
}

步骤2:实现请求合并中间件

核心思路是用一个线程安全的字典缓存正在处理的请求任务,相同请求(路径+查询+指定请求头)复用首个请求的处理结果,避免重复执行后端逻辑。

1. 完善标记特性(可选,支持自定义配置)

public class SquashableAttribute : Attribute
{
    // 配置:指定参与请求标识计算的请求头(比如Authorization、User-Agent等)
    public string[] IncludeHeaders { get; set; } = Array.Empty<string>();
}

2. 实现请求合并中间件

using System.Collections.Concurrent;
using System.Linq;
using System.Text;

public class SquasherMiddleware
{
    private readonly RequestDelegate _next;
    // 缓存正在处理的请求任务,键为请求唯一标识,值为响应结果的包装任务
    private readonly ConcurrentDictionary<string, Task<(int StatusCode, Dictionary<string, string> Headers, byte[] Body)>> _pendingRequests = new();

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

    public async Task InvokeAsync(HttpContext context)
    {
        var endpoint = context.GetEndpoint();
        if (endpoint == null)
        {
            await _next(context);
            return;
        }

        // 判断当前端点是否启用请求合并
        var squashableAttr = endpoint.Metadata.GetMetadata<SquashableAttribute>();
        if (squashableAttr == null)
        {
            await _next(context);
            return;
        }

        // 生成请求唯一标识:路径+查询字符串+指定请求头
        var requestKey = GenerateRequestKey(context, squashableAttr);

        // 尝试获取或创建处理任务:首个请求执行后端逻辑,后续请求复用结果
        var processingTask = _pendingRequests.GetOrAdd(requestKey, async key =>
        {
            try
            {
                // 捕获原始响应流,替换为内存流以读取响应内容
                var originalResponseBody = context.Response.Body;
                using var responseStream = new MemoryStream();
                context.Response.Body = responseStream;

                // 执行后续中间件及端点逻辑
                await _next(context);

                // 读取响应内容并保存
                responseStream.Seek(0, SeekOrigin.Begin);
                var bodyBytes = await responseStream.ToArrayAsync();
                var headers = context.Response.Headers.ToDictionary(kv => kv.Key, kv => kv.Value.ToString());
                
                return (context.Response.StatusCode, headers, bodyBytes);
            }
            finally
            {
                // 任务完成后从缓存移除,避免内存泄漏
                _pendingRequests.TryRemove(key, out _);
            }
        });

        try
        {
            // 等待处理任务完成,复用结果
            var (statusCode, headers, bodyBytes) = await processingTask;

            // 将复用的响应写入当前请求的Response
            context.Response.StatusCode = statusCode;
            foreach (var header in headers)
            {
                context.Response.Headers.TryAdd(header.Key, header.Value);
            }
            await context.Response.Body.WriteAsync(bodyBytes);
        }
        catch (Exception ex)
        {
            // 传播首个请求的异常给后续请求
            context.Response.StatusCode = StatusCodes.Status500InternalServerError;
            await context.Response.WriteAsync(ex.Message);
        }
    }

    private string GenerateRequestKey(HttpContext context, SquashableAttribute attr)
    {
        var sb = new StringBuilder();
        // 基础标识:路径+查询字符串
        sb.Append(context.Request.Path);
        sb.Append(context.Request.QueryString);

        // 添加指定的请求头,确保不同头的请求不被合并
        foreach (var headerName in attr.IncludeHeaders)
        {
            if (context.Request.Headers.TryGetValue(headerName, out var headerValue))
            {
                sb.Append($"|{headerName}:{headerValue}");
            }
        }

        return sb.ToString();
    }
}

public static class SquasherMiddlewareExtensions
{
    public static IApplicationBuilder UseSquasher(this IApplicationBuilder builder)
    {
        return builder.UseMiddleware<SquasherMiddleware>();
    }
}

3. 注册中间件

确保中间件注册在UseRouting()之后,UseEndpoints()之前:

.NET Core 3.1 Startup.cs

public void Configure(IApplicationBuilder app, IWebHostEnvironment env)
{
    // 其他中间件(如UseStaticFiles、UseAuthentication等)...
    app.UseRouting();
    // 注册请求合并中间件
    app.UseSquasher();
    app.UseEndpoints(endpoints =>
    {
        endpoints.MapControllers();
    });
}

.NET 6/7 Program.cs

var builder = WebApplication.CreateBuilder(args);
builder.Services.AddControllers();

var app = builder.Build();

// 其他中间件...
app.UseRouting();
app.UseSquasher();
app.MapControllers();

app.Run();

关键注意事项

  • 请求唯一标识:可根据业务需求调整GenerateRequestKey方法,比如添加用户ID(用户级请求隔离)、请求方法等,避免不同场景的请求被误合并。
  • 内存管理:任务完成后立即从字典移除,也可添加超时清理逻辑(比如用定时器定期清理超过5分钟的任务),防止内存泄漏。
  • 异常处理:确保首个请求的异常能正确传播给后续请求,避免静默失败。
  • 兼容性:代码兼容.NET Core 3.1、.NET 6和.NET 7,无需额外修改。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 00:45:36