首页
看点啥
插画图片
首页 科技看点 如何在 .NET 中优雅地合并并发请求实用指南

如何在 .NET 中优雅地合并并发请求实用指南

2026-08-31 0

平时做技术实践时,很多问题不是概念不会,而是细节没串起来。拿“如何在 .NET 中优雅地合并并发请求”来说,它看着像小点,放到项目里常会牵出环境、配置、兼容性和维护成本。下面按实际采用顺序,把思路、关键写法和容易踩坑的地方讲清楚,便于大家直接对照操作。

(链接已移除)

NuGet:RequestBatcher

理解这一步时,本文基于正式发布的 RequestBatcher v0.0.2 源码编写,兼容 .NET 8 及以上版本。

前言

从实现思路看,假设一个商品详情 API 会同时发查询价格。最直接的实现,是每个请求都单独访问一次数据库:

请求 A -> SELECT 商品 101
请求 B -> SELECT 商品 101
请求 C -> SELECT 商品 205

理解这一步时,请求 A 和 B 查询的是同一件商品,却仍然产生了两次数据库调用。流量集中到达时,数据库连接、网络往返、命令解析和查询执行都会被重复很多次。

实际处理时,要把同一时刻涌入的查询合起来,应用需一个能一次处理多项请求的批处理方法。它能够先去重商品 ID,再执行一次批量查询:

请求 A -> 商品 101 --\
请求 B -> 商品 101 ----> 批量处理:去重为 [101, 205] -> 一次批量查询
请求 C -> 商品 205 --/

实际处理时,这个批处理方法能够先对当前批次中的商品 ID 执行 Distinct,再借助一次批量查询读取数据,最后把结果分发给每个调用方。这种合同时只处理当前进程内短暂积压的请求,不是缓存,也不会为了凑批固定等待。这里的去重只发生在同一次批处理里;相同查询进入不同批次时,仍然会再次访问下游。

理解这一步时,批量写入也是同一类问题。对于某个数据写入场景,多个调用方各自提交一条数据记录,批处理方法能够把当前批次一次写入下游。

从实现思路看,当然,也能够让上游调用方先收集一个 List,再统一交给数据库。问题是,请求通常来自彼此独立的 HTTP 调用、消息处理器或后台任务。它们同时不知道同一时刻还有谁在做相同的事情,更不应该共同维护一套批次、并发、超时、异常和停止逻辑。

RequestBatcher 解决的正是这个协调问题:

从实现思路看,它采用两个会反复出现的概念:分区是独立、按顺序处理请求的内存队列;Consumer 是从一个分区取出当前请求同时调用批处理方法的处理循环。分区键只决定请求进入哪个分区,不负责去重。

  • 每个调用方仍然只提交自己的一个请求;
  • 结合项目来看,同一分区中已经排队的请求,会被合同时成最多 BatchSize 项的批次;
  • 批处理方法一次拿到一组请求,负责执行真正的下游批量操作;
  • 每个调用方仍然等待自己的 Task,同时得到这个请求实际的成功、失败或取消结果;
  • 队列容量和批处理同时发都能够被限制,避免流量突发直接压垮下游。

实际处理时,它不是持久化消息队列,也不是事务协调器。它更像应用调用路径中的一个进程内“请求汇合点”:上游保持单请求接口,下游获得批量处理机会。

公共接口

实际处理时,对应用来说,RequestBatcher 只有两个核心接口:IRequestBatcher 用来提交请求,IRequestBatchHandler 用来处理当前批次:

public interface IRequestBatcher
{
    Task ProcessAsync(
        TRequest request,
        CancellationToken cancellationToken = default);
    Task ProcessAsync(
        IEnumerable requests,
        CancellationToken cancellationToken = default);
}
public interface IRequestBatchHandler
{
    ValueTask HandleAsync(
        IReadOnlyList requests,
        CancellationToken cancellationToken = default);
}

理解这一步时,调用方通常采用第一个重载提交一项请求,同时等待它实际处理完成;第二个重载用来调用方本来就持有一组请求时的一次性提交。Handler 接收的是当前批次的 IReadOnlyList,而不是某次提交的原始集合;一次显式提交仍可能被拆成多个 Handler 调用。

批量处理的收益来源

在这个场景下,RequestBatcher 本身不会让一段逐条业务代码自动变快。收益来自 Handler 把多项请求真正变成更少的下游操作:

逐条处理:N 个请求 -> N 次数据库查询
批量处理:N 个请求 -> M 次 Handler -> M 次批量查询
                         通常 M < N

落到代码里,对于商品价格查询,示例 Handler 会先对当前批次的 ProductId 执行 Distinct,再用一次 ANY(@ProductIds) 查询价格,最后把结果写回每个 PriceQuery。所以,重复请求越集中在同一批次中,减少的数据库调用就越多。

从实现思路看,下面这组端到端基准已经提交到 PostgreSqlPriceQueryBenchmarks.cs。工作负载是 1,000 次查询,但只有 10 个商品 ID,每个重复 100 次;这是适合批内去重的特定场景,不代表一般查询的平均结果。

在这个场景下,直连路径和 RequestBatcher 都限制为 4 路同时发、Npgsql 连接池固定为 4 个连接。容器启动、建表、填充数据和连接池预热不计入耗时,计时范围是提交请求到全部 Task 完成。每个场景预热 2 次、测量 8 次,并校验每个请求都得到了正确价格;RequestBatcher 路径还会验证 Handler 确实形成了多项批次。

场景平均耗时标准差下游查询处理情况
逐条查询,1,000 个请求71.223 ms2.793 ms每个请求一次 SELECT1,000 次单项查询
落到代码里,RequestBatcher,10 个商品 ID、每个重复 1002.930 ms0.430 ms每个 Handler 批次一次 ANY(@product_ids)批内 Distinct 后查询,再逐项回填结果

理解这一步时,测试环境为 macOS 15.7.8、Apple M2 Max、.NET 8.0.25,以及 Testcontainers 启动的 PostgreSQL 17.6-alpineBatchSize=100MaxConcurrency=4。在这组高重复查询里,逐条查询的平均耗时约为 RequestBatcher 的 24.3 倍(71.223 / 2.930);RequestBatcher 约为直连路径的 4.1%。BenchmarkDotNet 也提示单次迭代不足 100 ms,所以这些数值只说明这个请求分布和这台机器上的相对结果,不应外推为所有查询场景的固定倍数。

从实现思路看,这不是“相同商品永远只查一次”的承诺。相同 ProductId 仍然可能落入不同 Handler 批次;它只会在当前批次中借助 Distinct 合同时。批次边界会受到请求到达时机、Consumer 调度、分区和 BatchSize 的共同影响。这组数据只说明:当大量重复查询确实在同一批次中汇合,并且 Handler 做了真正的批量查询时,收益会很明显。

适用场景

实际处理时,RequestBatcher 适合彼此独立、能够在当前进程内短暂排队,同时且下游更适合一次处理多项的请求。典型场景包括:

  • 理解这一步时,数据库、缓存或下游 API 兼容多项查询,希望把同时发单项查询合并,并在批内去重相同 Key;
  • 数据库兼容批量 INSERTUPDATEUPSERT,希望把同时发单条写入合并成较少的批量操作;
  • 在这个场景下,流量会在短时间内集中到达,需限制排队请求数和下游同时发量,对数据库或下游 API 形成并发限制与背压保护;
  • 在这个场景下,相关请求需保持分区内顺序,或者能够在当前批次中合同时重复更新、去重相同查询;
  • 理解这一步时,调用方取消时,只需撤销尚未分发的请求;已经交给 Handler 的请求应继续执行同时得到实际结果。

落到代码里,它尤其适合“上游接口必须保持单请求,只有下游知道如何批量处理”的调用链。调用方不需感知批次,Handler 也不需关心每个请求来自哪个调用方。

这里的限制由 MaxConcurrencyMaxPendingRequests 提供:前者限制同时执行的 Handler 数量,后者限制内存中尚未完成的请求数量。它们会把下游变慢的压力传回上游,但不按每秒请求数做固定速率限制;需 QPS 限流时,应与专门的限流器配合采用。

不适用场景

下面这些需求已经超出进程内请求合并的职责:

  • 已接收的请求必须在进程崩溃后恢复。此时需持久化存储或可靠消息队列;
  • 结合项目来看,操作必须与调用方当前事务一起提交或回滚。把操作移出原事务会改变一致性边界;
  • 落到代码里,调用方必须直接从一次调用中得到业务结果,同时且不能接受像下方示例那样由请求对象保存结果。RequestBatcher 只得到表示处理完成的 Task
  • 在这个场景下,下游操作依赖自动重试,或者副作用必须满足恰好一次(exactly-once)。RequestBatcher 不提供这些保证;
  • 实际处理时,业务必须凑满最小批量,或者依赖固定收集窗口。RequestBatcher 只处理当前已经排队的请求;
  • 实际处理时,单个调用方断开或取消后,正在执行的下游操作也必须立即停止。调用方的 CancellationToken 不会传给共享的 Handler 调用。

请求合并过程

RequestBatcher 的基本流程只有四步:

  • 调用方借助 ProcessAsync 提交一个 TRequest
  • 请求按设置进入某个内存分区,并获得独立的完成状态。
  • 落到代码里,Consumer 能够继续处理该分区时,取出当前已经排队的最多 BatchSize 个请求,只调用一次 Handler。
  • 结合项目来看,Handler 的执行结果再分别完成这一批请求对应的 Task

从实现思路看,这里最重要的词是“已经排队”。RequestBatcher 采用的是机会式合同时,而不是固定时间窗口。

在这个场景下,它不会在第一个请求到达后等待 10 ms、50 ms,试图凑满一个批次;Consumer 能够继续处理该分区时,会尽快取走当前已有的请求。所以,BatchSize 是单次 Handler 调用的上限,不是触发处理的最小数量。

一次可能的执行过程如下所示:

t0  请求 A 到达,分区空闲,Handler 开始处理 [A]
t1 A 仍在处理,请求 B、C、D 进入同一分区排队
t2 A 完成,Handler 下一次处理 [B, C, D]

实际处理时,具体批次边界取决于请求到达与 Consumer 调度的时机,RequestBatcher 不承诺 A 一定单独成批,也不承诺 B、C、D 一定在同一批。它只保证单批不超过 BatchSize,同时在有并发积压时有机会形成更大的批次。

这个取舍很实用:低流量下不为了凑批而主动增加延迟,高同时发下又能够借助排队请求减少下游调用次数。如果业务必须“至少凑够 100 条再处理”,或者必须采用固定收集窗口,RequestBatcher 并不适合。

更快接入示例

先安装 NuGet 包:

dotnet add package RequestBatcher --version 0.0.2

先定义查询请求和得到结果。ProcessAsync 只得到表示处理完成的 Task,所以请求对象需保存自己的查询结果:

public sealed record ProductPrice(
    long ProductId,
    decimal Price);
public sealed class PriceQuery(long productId)
{
    public long ProductId { get; } = productId;
    public ProductPrice? Result { get; private set; }
    public void SetResult(ProductPrice? result) => Result = result;
}
public interface IProductPriceStore
{
    Task> FindManyAsync(
        IReadOnlyList productIds,
        CancellationToken cancellationToken);
}

理解这一步时,Handler 对当前批次中的商品 ID 去重,执行一次批量查询,再把结果写回每个请求:

public sealed class PriceQueryHandler(IProductPriceStore store)
    : IRequestBatchHandler
{
    public async ValueTask HandleAsync(
        IReadOnlyList requests,
        CancellationToken cancellationToken = default)
    {
        var productIds = requests
            .Select(request => request.ProductId)
            .Distinct()
            .ToArray();
        var prices = await store.FindManyAsync(productIds, cancellationToken);
        var pricesByProductId = prices.ToDictionary(price => price.ProductId);
        foreach (var request in requests)
        {
            request.SetResult(
                pricesByProductId.GetValueOrDefault(request.ProductId));
        }
    }
}

随后把 Handler 和 RequestBatcher 一起注册到应用已有的 DI 容器:

builder.Services.AddRequestBatcher(
    ServiceLifetime.Scoped,
    options =>
    {
        options.BatchSize = 100;
        options.MaxConcurrency = 4;
        options.MaxPendingRequests = 10_000;
        options.FullMode = RequestBatchFullMode.Wait;
        // 可选:让相同商品路由到同一分区,增加它们在同一批次中去重的机会。
        // 分区键只决定路由,不会自动删除重复请求。
        options.UsePartitionKey(query => query.ProductId);
    });

最后,在应用服务中注入 IRequestBatcher。调用方仍然一次只查询一个商品:

public sealed class ProductPriceService(
    IRequestBatcher requestBatcher)
{
    public async Task GetAsync(
        long productId,
        CancellationToken cancellationToken = default)
    {
        var query = new PriceQuery(productId);
        await requestBatcher.ProcessAsync(query, cancellationToken);
        return query.Result;
    }
}

GetAsync 会等待查询所在的 Handler 批次完成,再读取这个请求对应的结果。相同 ProductId 借助分区键进入同一分区,所以有机会在当前批次中被去重;如果它们落入不同批次,仍然会再次访问数据库,RequestBatcher 不会替代缓存。

结合项目来看,若 Handler 没有需由 DI 解析的依赖,也能够借助 AddRequestBatcher 的委托重载直接注册一个 RequestBatchHandler,注册时仍需明确指定 Handler 生命周期。

Handler 生命周期

在这个场景下,RequestBatcher Coordinator 本身是 Singleton,但 Handler 生命周期由注册时显式指定:

  • Scoped:每个 Handler 批次新建一个异步 Scope,同时在批次结束后释放;
  • Transient:同样在每个批次的 Scope 中解析一次;
  • Singleton:所有批次复用同一个实例。

结合项目来看,数据库上下文等 Scoped 依赖适合采用 Scoped Handler。选择 SingletonMaxConcurrency > 1 时,多个分区可能同时调用同一个 Handler 实例,所以它必须是线程安全的。

请求完成与异常回传

从实现思路看,请求合同时最麻烦的部分,其实不是把多个对象放进一个数组,而是把批量执行结果重新对应回每个调用方。

从实现思路看,单个请求进入 RequestBatcher 后,会被包装成内部的 PendingBatchRequest。其中既保存请求本身,也保存状态、取消注册和完成源。

落到代码里,Handler 正常完成时,这一批请求对应的 Task 全部成功;Handler 抛出异常时,同一批调用方都会收到这个原始异常。失败不会触发自动重试,Consumer 会继续处理后续批次。

Caller A --\
Caller B ----> Handler([A, B, C]) 成功 ----> Task A/B/C 全部成功
Caller C --/
Caller D --\
Caller E ----> Handler([D, E]) 失败 ------> Task D/E 收到同一异常

Handler 得到 ValueTask,是因为同步完成的处理器能够避免额外新建 Task。这个 ValueTask 只会在 RequestBatcher 内部等待一次,不会直接暴露给调用方。

显式批量提交的语义

落到代码里,除了提交单个请求,RequestBatcher 也兼容提交调用方已经持有的一组请求:

await requestBatcher.ProcessAsync(priceQueries, cancellationToken);

这个重载只表示调用方用一次 ProcessAsync 提交了多个请求,不表示 Handler 只会调用一次。RequestBatcher 会先枚举输入同时新建快照,再让每一项独立路由;它们可能被 BatchSize 拆开,也可能进入不同分区并行执行。

行为单个请求显式提交一组请求
输入接收一个 TRequest枚举一次并保存快照;空序列立即完成
路由路由到一个分区每一项独立路由,一次提交能够跨分区
Handler 边界可能与其他排队请求合并可能按分区和 BatchSize 拆成多次 Handler 调用
完成Task 等待这一项一个 Task 等待组内所有项结束
失败得到所在 Handler 批次的异常任一相关项失败,整组 Task 失败;已经成功的操作不会回滚

显式提交的 Task 会等待所有项结束,再按下面的优先级决定最后结果:

  • 只要有任何一项失败,整组 Task 就失败;错误可能来自 Handler,也可能来自入队、路由或 Consumer;
  • 从实现思路看,没有失败但至少一项取消,整组 Task 取消;
  • 所有项都成功,整组 Task 成功。

落到代码里,一次提交可能跨多个 Handler 批次。如果多个批次分别失败,所有不同的异常实例仍能够从 Task.Exception.InnerExceptions 中取得;普通 await 则遵循 .NET 的 Task 语义抛出其中一个异常。

容量策略取决于 FullModeWait 能够随着容量释放逐步接收超大请求组,Fail 则要求整组立即获得容量。具体行为见下文的“背压与容量控制”。无论哪种模式,一次显式提交都不是事务边界、分区边界或 Handler 批次边界。

批次大小与处理并发

BatchSize 控制一次 Handler 最多处理多少项,MaxConcurrency 控制最多有多少次 Handler 调用同时行执行。

它们不能互相替代:批次很大不代表同时发很高,并发很高也不代表每批有很多请求。

BatchSize:单批上限

假设 BatchSize = 100

  • 分区中只有 3 项排队时,Handler 能够收到 3 项;
  • 分区中有 250 项排队时,会被拆成最多 100 项的多个批次;
  • 不同分区的请求不会被放进同一个 Handler 批次。

结合项目来看,批次大小应该根据下游能力决定,而不是越大越好。SQL 参数数量、请求体大小、事务持锁时间、缓存 Pipeline 长度和 API 限流,都可能成为新的边界。

MaxConcurrency:并发与分区

结合项目来看,RequestBatcher 采用 MaxConcurrency 个内存分区,同时为每个分区建立一个顺序消费循环。因此,这个设置同时决定 Handler 最大并发数和处理分区数。

设置路由方式顺序语义
MaxConcurrency = 1所有请求进入唯一分区保持全局处理顺序
MaxConcurrency > 1,无分区键请求逐项轮询到不同分区只保证各分区内部顺序
MaxConcurrency > 1,有分区键相同键路由到同一分区相同键进入分区后顺序处理,不同分区能够并行

落到代码里,同时发调用方在请求真正进入分区之前没有额外的先后保证。提高 MaxConcurrency 后,也不再存在跨分区的全局 FIFO。

实际处理时,下游如果最多只允许 8 个同时发连接,把 MaxConcurrency 配成 100 并不会带来免费性能,反而可能把压力转移到连接池和下游限流队列。它应该与下游的真实并发能力一起设置。

分区键与去重

更快接入中的价格查询采用 ProductId 作为分区键,目的是让相同商品的查询进入同一分区。真正的去重仍然来自 Handler 中的 Distinct;分区键本身不会删除任何请求。

落到代码里,分区键也能够用来需顺序处理的写入。比如,同一个商品价格的多个版本能够批量更新,但不希望两个 Handler 同时更新同一商品。

这时能够设置分区键:

builder.Services.AddRequestBatcher(
    ServiceLifetime.Scoped,
    options =>
    {
        options.BatchSize = 100;
        options.MaxConcurrency = 4;
        options.MaxPendingRequests = 10_000;
        options.UsePartitionKey(update => update.ProductId);
    });

相同 ProductId 会进入同一分区,因而不会被两个 Handler 调用同时发处理;其他商品仍然能够在不同分区并行执行。

Handler 还能够合同时当前批次中的重复更新,只保留版本最高的一项:

public sealed record PriceUpdate(
    long ProductId,
    long Version,
    decimal Price);
public sealed class PriceUpdateHandler(ProductPriceStore store)
    : IRequestBatchHandler
{
    public async ValueTask HandleAsync(
        IReadOnlyList requests,
        CancellationToken cancellationToken = default)
    {
        var latestUpdates = requests
            .GroupBy(request => request.ProductId)
            .Select(group => group.MaxBy(request => request.Version)!)
            .ToArray();
        await store.UpsertLatestAsync(latestUpdates, cancellationToken);
    }
}

但需留意,分区键只负责路由:

  • 它不会把同一个键的所有请求永久收集到同一个批次;
  • 它不会自动删除重复请求;
  • 不同键可能映射到同一个分区;
  • 同一键在当前批次中合并了,也可能在后续批次再次出现。

在这个场景下,因此,跨批次正确性仍然要由存储层保证。比如,价格表的 UPSERT 应只允许更高版本覆盖旧版本。分区键避免同一商品被同时发处理,数据库版本条件则避免后到的旧数据覆盖新状态,两者解决的不是同一个问题。

在这个场景下,数值型分区键必须是有限整数,字符串键不能为 null。键选择器应该稳定、无副作用、能够被同时发调用,并且只表达业务上的顺序需求。

在这个场景下,若显式组提交中的 Partition Key 选择器抛出异常,整组会失败且不会只处理其中一部分,已经预留的容量也会被释放。

批量写入场景

实际处理时,查询是本文的主例子,但相同的协调方式也能够用来批量写入。对于某个数据写入场景,不同调用方逐条提交数据,Handler 再将当前批次一次写入下游:

public sealed record DataWriteRequest(
    long DataId,
    string Payload);
public interface IDataStore
{
    Task WriteBatchAsync(
        IReadOnlyList requests,
        CancellationToken cancellationToken);
}
public sealed class DataWriteBatchHandler(IDataStore store)
    : IRequestBatchHandler
{
    public async ValueTask HandleAsync(
        IReadOnlyList requests,
        CancellationToken cancellationToken = default)
    {
        await store.WriteBatchAsync(requests, cancellationToken);
    }
}

注册和调用方式与查询相同:

builder.Services.AddRequestBatcher(
    ServiceLifetime.Scoped,
    options =>
    {
        options.BatchSize = 256;
        options.MaxConcurrency = 4;
        options.MaxPendingRequests = 10_000;
        options.FullMode = RequestBatchFullMode.Wait;
    });
public sealed class DataWriteService(
    IRequestBatcher requestBatcher)
{
    public Task WriteAsync(
        DataWriteRequest request,
        CancellationToken cancellationToken = default) =>
        requestBatcher.ProcessAsync(request, cancellationToken);
}

WriteAsync 得到的 Task 表示这条数据的实际处理结果,而不是仅表示它已经进入内存队列。真正的批量写入应在 Handler 中完成,比如采用数据库数组参数、批量写入 Pipeline 或下游 API 的批量接口。

实际处理时,项目仓库中的 PostgreSQL 示例同时展示了批量 UPSERT、查询去重和结果回填。

背压与容量控制

BatchSize 只能限制一次 Handler 的输入数量,不能限制总共有多少请求正在等待和执行。

结合项目来看,若生产速度长期高于消费速度,无界排队最后只会把下游过载变成应用内存过载。所以 RequestBatcher 还提供 MaxPendingRequests,限制已经接收、仍在排队或正在由 Handler 处理的请求数量;Wait 模式下,一次显式提交能够超过这个数量。

默认设置如下所示:

选项默认值含义
BatchSize128单次 Handler 调用的请求上限
MaxConcurrency1Handler 最大并发数,同时也是分区数
MaxPendingRequests8192已接收且尚未完成的请求上限;Wait 模式下显式提交可超过该值
FullModeWait落到代码里,容量不足时异步等待;也能够选择 Fail
UsePartitionKey(...)未设置默认逐项轮询路由

Wait 模式

RequestBatchFullMode.Wait 会异步等待容量。调用方不会阻塞线程,同时且在等待期间能够借助自己的 CancellationToken 取消。

结合项目来看,这种模式把背压自然传回调用链:下游变慢后,上游的 ProcessAsync 也会变慢,而不是继续无限制接收请求。

ProcessAsync(IEnumerable) 而言,超大请求组也会被处理。假设 MaxPendingRequests = 10_000,一次提交 20_000 项:前面最多 10_000 项会先进入队列,剩余请求随着前面请求完成、容量释放而继续进入。调用方仍只等待一个 Task,它会在整组请求都成功、失败或取消后结束。这是逐步入队的策略,不是事务或 Handler 批次。

Fail 模式

RequestBatchFullMode.Fail 在容量不足时立即让得到的 TaskRequestBatchQueueFullException 失败:

try
{
    await requestBatcher.ProcessAsync(request, cancellationToken);
}
catch (RequestBatchQueueFullException exception)
{
    logger.LogWarning(
        "Request batch queue is full. Capacity: {Capacity}, Requested: {Requested}",
        exception.Capacity,
        exception.RequestedCount);
    // 按业务约定返回 429、降级或交给可靠队列
}

Fail 模式下,显式提交必须整组获得当前可用容量;不会出现前半组已经入队、后半组因为容量不足被拒绝的状态。整组数量超过 MaxPendingRequests,或当前容量不足时,都会以 RequestBatchQueueFullException 失败。

请求取消的时机

理解这一步时,一个 Handler 批次可能同时包含多个调用方的请求。如果把其中任意一个调用方的 CancellationToken 直接传给 Handler,那么一个 HTTP 客户端断开连接,就可能把其他调用方共享的数据库操作一起取消。

理解这一步时,更麻烦的是,Handler 开始执行后可能已经产生了副作用。此时把调用方的 Task 标记为取消,会让上游误以为操作没有发生,重试后反而造成重复写入。

RequestBatcher 所以只在 Handler 分发之前接受调用方取消:

Queued ---------> Processing ---------> Succeeded / Faulted
  | |
  | caller cancel | caller cancel
  v v
Canceled 继续等待 Handler 的真实结果

具体语义是:

  • 落到代码里,请求还在等待容量或排队时,调用方取消能够移除这项请求同时取消它的 Task
  • 从实现思路看,Consumer 准备处理时,会用原子状态切换把请求从 Queued 改为 Processing
  • 落到代码里,一旦切换成功,之后的调用方取消不再改变结果;Task 最后反映 Handler 的成功或异常。
  • 实际处理时,传给 Handler 的 Token 属于 RequestBatcher 的 Consumer 生命周期,不是任意一个调用方的 Token。

实际处理时,这不是忽略取消,而是避免取消状态掩盖一个可能已经发生的共享副作用。即使 BatchSize = 1,RequestBatcher 也不会把调用方 Token 传给 Handler。业务如果要求“调用方一断开,正在执行的下游操作必须立刻停止”,RequestBatcher 就不适合这条调用路径,应直接调用下游,或采用另一套与调用方生命周期绑定的取消机制。

应用关闭时的请求处理

RequestBatchCoordinator 实现了 IAsyncDisposable,同时提供 StopAsync。停止过程按下面的顺序执行:

停止接收新请求
    -> 处理完停止前已经开始的所有提交,包括仍在等待容量的部分
    -> 停止 Consumer
    -> 释放生命周期资源

开始关闭后提交的新请求会以 ObjectDisposedException 失败;在关闭开始前已经提交、但仍在等待容量的请求会继续等待同时处理完成。调用方自己的 Token 更早取消时,仍可能观察到取消结果。在 Consumer 正常运行的前提下,已经接收的请求会继续处理到结束。

传给 StopAsyncCancellationToken 只取消调用方对停止过程的等待,不会撤销已经开始的处理。即使等待 StopAsync(token) 时传入的 Token 已取消,后台停止过程仍会继续。

实际处理时,这类生命周期语义很重要。请求合同时通常处在 HTTP、数据库和应用宿主之间,如果应用停止时直接取消所有 Consumer,调用方可能永远等不到结果,已经接收的业务请求也会无声丢失。

内部实现:从入队到完成

理解这一步时,RequestBatcher 的公共 API 很小,真正的协调由几个内部组件完成:

理解这一步时,RequestBatcher 底层采用我开源的 BufferQueue。

IRequestBatcher.ProcessAsync
    |
    v
RequestBatchCoordinator
    |
    +-- PendingRequestProducer
    | |
    | v
    | BufferQueue Memory Topic
    | |
    | +-- Partition 0 -> Consumer 0
    | +-- Partition 1 -> Consumer 1
    | +-- ...
    |
    +-- PendingBatchRequest
            |
            +-- request
            +-- queued / processing / canceled / completed
            +-- completion

注册 AddRequestBatcher 时,内部会新建一个 BufferQueue Memory Topic:

  • PartitionNumber 采用 MaxConcurrency
  • BoundedCapacity 采用 MaxPendingRequests
  • 理解这一步时,队列满策略映射自 RequestBatcher 的 WaitFail
  • 落到代码里,设置了 Partition Key 时,路由规则一起交给内部 Topic。

落到代码里,Coordinator 随后新建与 MaxConcurrency 相同数量的 Pull Consumer,也就是主动从分区取请求的处理循环。每个 Consumer 顺序读取自己的分区,每次最多拉取 BatchSize 项,同时关闭自动提交(Auto Commit):Handler 处理结束后,再由 Consumer 显式提交这一批的消费进度。

实际处理时,Consumer 拿到一批内部请求后,会先跳过已经取消的项,再把剩余的 TRequest 复制到数组中交给 Handler。Handler 正常得到或抛出异常时,Consumer 都会先完成对应请求的状态,再提交这一批的消费进度。读取批次、Consumer 循环或提交进度本身失败时,属于消费基础设施错误;这时不会保证本批进度已提交。

落到代码里,内部监控组件(Monitor)会记录第一个 Consumer 故障,停止继续接收请求,同时让仍在排队的请求以该异常失败,避免调用方持有一个永远不会完成的 Task

理解这一步时,显式提交一组请求时,内部会新建共享的 BatchSubmissionCompletion。它只保留一个 TaskCompletionSource:组内每项结束时更新共享计数和异常状态,最后一项结束时再完成这个 Task。这样不需为组内每项额外新建完成源,也不需构造 Task[] 后调用 Task.WhenAll;它只优化完成聚合的分配与协调开销,不改变取消、异常和 Handler 批次语义。

在这个场景下,BufferQueue 在这里是实现细节。应用不需注册 Topic、Producer 或 Consumer,也不应该依赖内部 Topic 名称;对外契约始终只有 IRequestBatcherIRequestBatchHandler 和设置项。

采用限制

确定场景适合之后,还需确认几项实现上的限制:

  • ProcessAsync(IEnumerable) 表示一次提交,不保证只产生一次 Handler 调用。
  • BatchSize 是单批上限,不是最小数量;低流量下可能持续产生单项批次。
  • MaxConcurrency > 1 时只保证分区内顺序,不保证跨分区全局 FIFO。
  • 结合项目来看,Partition Key 只负责路由,不会自动去重,也不保证一个 Key 独占一个分区或进入同一批次。
  • 显式组提交跨批次失败时,不会回滚已经成功的项;它不是事务边界。
  • 理解这一步时,Singleton Handler 在同时发大于 1 时必须线程安全;Scoped 或 Transient Handler 则按处理批次新建 Scope。
  • MaxPendingRequests 只限制内存中的未完成请求,同时形成背压;它不提供可靠投递保证。

总结

实际处理时,优雅地合同时并发请求,不是让每个调用方都学会收集批次,而是把“单项提交”和“批量执行”分成两个独立契约。

落到代码里,RequestBatcher 让调用方继续采用轻松的 ProcessAsync(request),由内部完成机会式合同时、分区路由、并发限制、背压、结果回传和有序停止;Handler 只关心如何把 IReadOnlyList 变成一次真正的数据库、缓存或下游批量操作。

选择它之前,能够先确认三件事:

  • 请求允许在当前进程内短暂排队,进程失败后不要求自动恢复;
  • 下游确实有批量处理能力,并且批量能够减少固定成本;
  • 实际处理时,业务接受“BatchSize 是上限、调用方取消只在分发前有效、分区内有序”的处理语义。

实际处理时,满足这些条件时,它能够把原本分散在各个调用方中的批次协调逻辑收回到一个清晰边界里,让上游保持轻松,也让下游真正获得批量处理的机会。

项目地址:eventhorizon-cli/RequestBatcher

本文对应版本:RequestBatcher v0.0.2

完整示例:RequestBatcher.Deduplication

中文文档:RequestBatcher README

底层队列:eventhorizon-cli/BufferQueue

核心实现:RequestBatchCoordinator | RequestBatchConsumer

行为测试:RequestBatchCoordinatorTests | RequestBatchSubmissionTests

NuGet:RequestBatcher

到此这篇关于在 .NET 中优雅地合同时并发请求的文章就介绍到这了,更多相关在 .NET 中优雅地合并并发请求内容请搜索脚本之家以前的文章或继续浏览下面的相关文章希望大家以后多多兼容脚本之家!

您可能感兴趣的文章:
  • 如何在 .NET 中优雅地合并并发请求
  • asp.net借助消息队列处理高同时发请求(以抢小米手机为例)
  • .net core同时发请求发送HttpWebRequest的坑解决
  • 让Win2008+IIS7+ASP.NET兼容10万同时发请求
喜欢(0)

上一篇

MariaDB数据库的外键约束实例完整指南

下一篇

excel vba 高亮显示当前行代码实用指南

猜你喜欢