Voocii博客
首页博客AI 热榜作品集读书友链工具关于

© 2026 Voocii. Built with Next.js & tRPC.

GitHubXEmailRSS


一文搞明白.NET 中的 Pipeline

dotnetrick-hayekrick-hayek2025年3月30日

引言

"Pipeline"(管道)是 .NET 生态里一个反复出现的词:ASP.NET Core 用它处理 HTTP 请求,HttpClient 用它组织出站消息,System.IO.Pipelines 用它做高性能字节流处理,Channel<T> 和 TPL Dataflow 用它做生产者-消费者流水线。它们名字相似,但解决的问题、内部机制完全不同。这里简单介绍一下这几种 Pipeline :机制原理、请求/数据的具体流转过程、以及对应的示例代码。

一、Pipeline 的本质:责任链模式

抛开具体实现,Pipeline 在设计模式层面对应的是责任链模式(Chain of Responsibility):把一个复杂的处理过程拆成一串独立的处理单元(节点),每个节点只关心自己的逻辑,处理完后把请求/数据交给下一个节点,直到链的末端。

它的价值在于:

  • 解耦:每个节点不需要知道整条链的全貌,只需要知道"我处理完之后交给谁"。
  • 可组合:节点的顺序、数量都可以在运行时动态调整。
  • 可短路:任何一个节点都可以决定"到此为止,不再往下传"。

.NET 里的各种 Pipeline,都是这个思想在不同场景下的具体化——只是"节点"和"传递的数据"不一样:ASP.NET Core 里节点是中间件、传递的是 HttpContext;System.IO.Pipelines 里节点是读写双方、传递的是内存块。

下图展示了责任链模式的抽象结构,也是本文所有具体 Pipeline 的共同骨架:

责任链模式抽象结构

二、ASP.NET Core 中间件管道

2.1 核心概念

ASP.NET Core 的请求处理管道由一系列 中间件(Middleware) 组成,每个中间件都是一个 RequestDelegate:

public delegate Task RequestDelegate(HttpContext context);

在 Program.cs(或早期版本的 Startup.Configure)里,我们通过 IApplicationBuilder 把中间件一个个"注册"进管道:

var app = builder.Build();

app.Use(async (context, next) =>
{
    Console.WriteLine("中间件 A:进入");
    await next(context);           // 调用下一个中间件
    Console.WriteLine("中间件 A:返回");
});

app.Use(async (context, next) =>
{
    Console.WriteLine("中间件 B:进入");
    await next(context);
    Console.WriteLine("中间件 B:返回");
});

app.Run(async context =>
{
    // 管道的终端节点,不再调用 next
    await context.Response.WriteAsync("Hello Pipeline");
});

app.Run();

IApplicationBuilder.Build() 实际做的事情,是把这一串委托从后往前嵌套包裹成一个大的 RequestDelegate,类似:

RequestDelegate pipeline = ctx => middlewareA(ctx, ctx2 => middlewareB(ctx2, terminal));

2.2 洋葱模型与请求流转

这种"层层包裹"的结构,使得请求进入和响应返回会经过同一组中间件的两次——这就是著名的洋葱模型(Onion Model):请求像穿过洋葱一样从外层往里走,到达核心(终端处理器)后再从里往外穿出来。

具体流转顺序如下:

  1. 请求进入最外层中间件(比如异常处理中间件);
  2. 调用 next(),进入下一层(比如路由中间件);
  3. 依次深入,直到最内层的终端处理器(Endpoint / app.Run)生成响应;
  4. 响应沿着调用栈原路返回,依次经过每个中间件 next() 之后的代码;
  5. 最终响应到达客户端。

下图是这条洋葱模型的可视化:

ASP.NET Core 中间件洋葱模型

2.3 短路(Short-circuiting)

任何中间件都可以选择不调用 next(),直接终止管道并返回响应——这在鉴权、限流、缓存命中等场景非常常见:

app.Use(async (context, next) =>
{
    if (!context.Request.Headers.ContainsKey("X-Api-Key"))
    {
        context.Response.StatusCode = 401;
        await context.Response.WriteAsync("Missing API Key");
        return; // 不调用 next,管道到此短路
    }
    await next(context);
});

2.4 Use / Run / Map 的区别

方法作用
app.Use注册一个可以继续调用 next 的中间件,是链条的"中间节点"
app.Run注册一个终端中间件,不接受 next,是链条的"末端"
app.Map / app.MapWhen按路径或条件把请求分流到一条独立的子管道

三、HttpClient 的消息处理管道(DelegatingHandler)

如果说 ASP.NET Core 中间件管道处理的是"入站"请求,那么 HttpClient 的处理器链处理的就是"出站"请求。

3.1 结构

HttpClient 本身不发请求,真正发请求的是 HttpMessageHandler。默认的最终处理器是 HttpClientHandler(负责真正的 socket 通信),在它之前可以插入任意多个 DelegatingHandler,组成一条链:

HttpClient → Handler1 → Handler2 → ... → HttpClientHandler → 网络

3.2 自定义 DelegatingHandler

public class LoggingHandler : DelegatingHandler
{
    protected override async Task<HttpResponseMessage> SendAsync(
        HttpRequestMessage request, CancellationToken cancellationToken)
    {
        Console.WriteLine($"请求: {request.Method} {request.RequestUri}");

        var response = await base.SendAsync(request, cancellationToken); // 调用下一个 Handler

        Console.WriteLine($"响应: {(int)response.StatusCode}");
        return response;
    }
}

public class RetryHandler : DelegatingHandler
{
    protected override async Task<HttpResponseMessage> SendAsync(
        HttpRequestMessage request, CancellationToken cancellationToken)
    {
        for (int attempt = 1; attempt <= 3; attempt++)
        {
            var response = await base.SendAsync(request, cancellationToken);
            if (response.IsSuccessStatusCode || attempt == 3)
                return response;

            await Task.Delay(200 * attempt, cancellationToken);
        }
        throw new InvalidOperationException("不可达");
    }
}

在 IHttpClientFactory 中注册,顺序决定了链条的先后:

builder.Services.AddHttpClient("MyApi")
    .AddHttpMessageHandler<LoggingHandler>()
    .AddHttpMessageHandler<RetryHandler>();

这里同样是洋葱模型:请求先经过 LoggingHandler,再到 RetryHandler,最后到底层 HttpClientHandler 发出,响应再原路返回。

四、System.IO.Pipelines:高性能字节流水线

4.1 为什么需要它

传统基于 Stream 的读写在处理 TCP 这类协议时有几个痛点:

  • 数据边界不确定:一次 Read 可能读到半条消息("半包"),也可能读到多条消息("粘包"),需要自己维护缓冲区拼接。
  • 内存分配和拷贝多:为了拼接不完整的数据,经常需要 byte[] 之间来回拷贝,GC 压力大。
  • 背压(backpressure)难处理:生产者写得快、消费者读得慢时,容易内存暴涨。

System.IO.Pipelines(PipeReader / PipeWriter)就是为了解决这些问题而设计的高性能 API,核心思想是用环形链表的内存块(Segment)+ 引用计数代替一次性大数组,读写双方通过同一个 Pipe 对象协作。

4.2 核心流转过程

  1. 写入方通过 PipeWriter.GetMemory() 拿到一块可写内存,写完数据后调用 Advance(n) 告诉 Pipe "我写了 n 个字节",再调用 FlushAsync() 把数据提交给读取方,如果读取方处理太慢,FlushAsync 会根据背压配置异步等待;
  2. 读取方通过 PipeReader.ReadAsync() 拿到一个 ReadResult,里面的 ReadOnlySequence<byte> 可能横跨多个内存块;
  3. 读取方在 ReadOnlySequence<byte> 里查找完整的消息边界(比如换行符),如果找到就处理,如果没找到就调用 AdvanceTo(consumed, examined),其中 examined 之前的数据"已检查但未消费",Pipe 会保留这部分等待更多数据到达;
  4. 已经被消费的内存块会被 Pipe 回收复用,减少 GC 压力。

下图展示了 PipeWriter 与 PipeReader 之间基于内存段的协作流程:

System.IO.Pipelines 数据流转

4.3 示例代码:按行解析 TCP 数据

async Task ProcessLinesAsync(PipeReader reader, CancellationToken ct)
{
    while (true)
    {
        ReadResult result = await reader.ReadAsync(ct);
        ReadOnlySequence<byte> buffer = result.Buffer;

        while (TryReadLine(ref buffer, out ReadOnlySequence<byte> line))
        {
            ProcessLine(line); // 处理一条完整的消息
        }

        // 告诉 Pipe:buffer.Start 之前的数据已消费,
        // buffer.End 之前的数据已检查(没有更多完整行了)
        reader.AdvanceTo(buffer.Start, buffer.End);

        if (result.IsCompleted)
            break;
    }

    await reader.CompleteAsync();
}

bool TryReadLine(ref ReadOnlySequence<byte> buffer, out ReadOnlySequence<byte> line)
{
    var position = buffer.PositionOf((byte)'\n');
    if (position == null)
    {
        line = default;
        return false;
    }

    line = buffer.Slice(0, position.Value);
    buffer = buffer.Slice(buffer.GetPosition(1, position.Value)); // 跳过 \n
    return true;
}

五、Channel<T>:生产者-消费者管道

System.Threading.Channels 提供了一个线程安全、支持背压的队列,天然适合搭建"生产者写、消费者读"的流水线,是 BlockingCollection 的异步升级版。

var channel = Channel.CreateBounded<int>(capacity: 100); // 有界,触发背压

// 生产者
_ = Task.Run(async () =>
{
    for (int i = 0; i < 1000; i++)
    {
        await channel.Writer.WriteAsync(i); // 队列满时会异步等待
    }
    channel.Writer.Complete();
});

// 消费者
await foreach (var item in channel.Reader.ReadAllAsync())
{
    Console.WriteLine($"处理: {item}");
}

支持多生产者、多消费者(SingleReader/SingleWriter 选项可以做性能优化),常用于日志批处理、任务队列等场景。

六、TPL (Task Parallel Library) Dataflow:可编排的多阶段管道

如果流水线有多个处理阶段,且每个阶段的并行度、缓冲策略都不同,System.Threading.Tasks.Dataflow(TPL Dataflow)提供了现成的构建块:

var download = new TransformBlock<string, string>(async url =>
{
    using var http = new HttpClient();
    return await http.GetStringAsync(url);
}, new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 4 });

var parse = new TransformBlock<string, int>(html => html.Length);

var save = new ActionBlock<int>(len =>
{
    Console.WriteLine($"长度: {len}");
});

// 用 LinkTo 把各阶段串成管道
var linkOptions = new DataflowLinkOptions { PropagateCompletion = true };
download.LinkTo(parse, linkOptions);
parse.LinkTo(save, linkOptions);

download.Post("https://example.com");
download.Complete();
await save.Completion;

BufferBlock / TransformBlock / ActionBlock 之间通过 LinkTo 连接,每个 Block 内部自带缓冲队列和并行度控制,非常适合搭建"下载 → 解析 → 落库"这类多阶段异步流水线。

七、常见 Pipeline 一览

Pipeline所在层节点类型典型场景
ASP.NET Core 中间件管道Web 框架Middleware(RequestDelegate)鉴权、日志、路由、异常处理
HttpClient 处理器链出站 HTTPDelegatingHandler重试、日志、熔断、签名
System.IO.Pipelines底层 I/OPipeReader / PipeWriter自定义协议解析、高性能网络服务
Channel<T>并发编程Writer / Reader生产者-消费者、任务队列
TPL Dataflow并发编程TransformBlock / ActionBlock 等多阶段并行数据处理
IAsyncEnumerable + LINQ语言/运行时迭代器流式数据的惰性处理

八、如何选择

  • 处理 HTTP 请求/响应、需要在框架层面插入横切逻辑 → ASP.NET Core 中间件;
  • 需要给 出站请求 加统一的重试/日志/认证 → DelegatingHandler;
  • 需要解析 自定义二进制/文本协议,且对内存分配和吞吐极度敏感 → System.IO.Pipelines;
  • 简单的 生产者-消费者 场景,逻辑单一 → Channel<T>;
  • 多阶段、多并行度 的复杂数据流水线 → TPL Dataflow。

九、在真实系统里它们是怎么组合的

这些 Pipeline 并不是“二选一”的关系,而是不同层次的协作。一个典型的 Web 服务,通常会同时用到几种,但各自负责不同阶段。

例 1:一个订单服务 / API 服务

  • 进入 Web 服务时,ASP.NET Core 中间件负责鉴权、日志、异常处理、限流;
  • 当服务需要调用支付、库存、通知服务时,用 HttpClient 的 DelegatingHandler 做统一重试、熔断、签名;
  • 对于高并发下的订单创建/消息通知,使用 Channel<T> 做缓冲队列,把“收到请求”与“后续异步处理”解耦;
  • 如果订单处理要经历“校验 → 扣库存 → 生成发票 → 推送消息”几个阶段,可再用 TPL Dataflow 把这几个阶段切成并行/串行的 Block;
  • 如果服务本身接收自定义 TCP/UDP/二进制协议,或者做高性能网关,就会用 System.IO.Pipelines 处理字节流。

可以把它理解成:

“请求进入 → 中间件做框架层横切 → HttpClient 做出站调用 → Channel<T> / TPL Dataflow 做异步工作流”。

例 2:高吞吐日志采集 / 消息网关

  • 系统底层用 System.IO.Pipelines 解析网络字节流;
  • 解析后的消息进入 Channel<T> 做缓冲和背压控制;
  • 再由 TPL Dataflow 做多阶段处理(过滤、聚合、落盘、推送 Kafka/Redis);
  • 如果是 Web 接口暴露,也会加 ASP.NET Core 中间件做鉴权和限流。

例 3:实时数据处理服务

  • Web 入口用中间件做鉴权/日志;
  • 通过 DelegatingHandler 调用上游服务;
  • 使用 Channel<T> 接收实时事件;
  • 用 TPL Dataflow 做多阶段计算;
  • 如果需要接入自定义协议或极致性能,也会在底层引入 System.IO.Pipelines。

一句话总结

  • 中间件解决“HTTP 入口的统一处理”;
  • DelegatingHandler 解决“出站请求的统一处理”;
  • System.IO.Pipelines 解决“字节流/协议层的高性能处理”;
  • Channel<T> / TPL Dataflow 解决“应用内异步、并行、复杂工作流”。

结语

.NET 里的各种 Pipeline 表面上互不相干,但本质都是责任链模式的变体:把一个大任务拆成若干个可以独立开发、独立测试、按需编排的小节点。理解了这套"进入 → 处理/传递 → 短路或返回"的通用流转逻辑,再看具体的中间件、消息处理器、PipeReader/PipeWriter,就只是"节点长什么样、数据是什么类型"的差异了。

评论 (0)

暂无评论,快来抢沙发吧!

目录
  • 引言
  • 一、Pipeline 的本质:责任链模式
  • 二、ASP.NET Core 中间件管道
  • 三、HttpClient 的消息处理管道(DelegatingHandler)
  • 四、System.IO.Pipelines:高性能字节流水线
  • 五、Channel<T>:生产者-消费者管道
  • 六、TPL (Task Parallel Library) Dataflow:可编排的多阶段管道
  • 七、常见 Pipeline 一览
  • 八、如何选择
  • 九、在真实系统里它们是怎么组合的
  • 结语