"Pipeline"(管道)是 .NET 生态里一个反复出现的词:ASP.NET Core 用它处理 HTTP 请求,HttpClient 用它组织出站消息,System.IO.Pipelines 用它做高性能字节流处理,Channel<T> 和 TPL Dataflow 用它做生产者-消费者流水线。它们名字相似,但解决的问题、内部机制完全不同。这里简单介绍一下这几种 Pipeline :机制原理、请求/数据的具体流转过程、以及对应的示例代码。
抛开具体实现,Pipeline 在设计模式层面对应的是责任链模式(Chain of Responsibility):把一个复杂的处理过程拆成一串独立的处理单元(节点),每个节点只关心自己的逻辑,处理完后把请求/数据交给下一个节点,直到链的末端。
它的价值在于:
.NET 里的各种 Pipeline,都是这个思想在不同场景下的具体化——只是"节点"和"传递的数据"不一样:ASP.NET Core 里节点是中间件、传递的是 HttpContext;System.IO.Pipelines 里节点是读写双方、传递的是内存块。
下图展示了责任链模式的抽象结构,也是本文所有具体 Pipeline 的共同骨架:
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));
这种"层层包裹"的结构,使得请求进入和响应返回会经过同一组中间件的两次——这就是著名的洋葱模型(Onion Model):请求像穿过洋葱一样从外层往里走,到达核心(终端处理器)后再从里往外穿出来。
具体流转顺序如下:
next(),进入下一层(比如路由中间件);app.Run)生成响应;next() 之后的代码;下图是这条洋葱模型的可视化:
任何中间件都可以选择不调用 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);
});
Use / Run / Map 的区别| 方法 | 作用 |
|---|---|
app.Use | 注册一个可以继续调用 next 的中间件,是链条的"中间节点" |
app.Run | 注册一个终端中间件,不接受 next,是链条的"末端" |
app.Map / app.MapWhen | 按路径或条件把请求分流到一条独立的子管道 |
如果说 ASP.NET Core 中间件管道处理的是"入站"请求,那么 HttpClient 的处理器链处理的就是"出站"请求。
HttpClient 本身不发请求,真正发请求的是 HttpMessageHandler。默认的最终处理器是 HttpClientHandler(负责真正的 socket 通信),在它之前可以插入任意多个 DelegatingHandler,组成一条链:
HttpClient → Handler1 → Handler2 → ... → HttpClientHandler → 网络
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:高性能字节流水线传统基于 Stream 的读写在处理 TCP 这类协议时有几个痛点:
Read 可能读到半条消息("半包"),也可能读到多条消息("粘包"),需要自己维护缓冲区拼接。byte[] 之间来回拷贝,GC 压力大。System.IO.Pipelines(PipeReader / PipeWriter)就是为了解决这些问题而设计的高性能 API,核心思想是用环形链表的内存块(Segment)+ 引用计数代替一次性大数组,读写双方通过同一个 Pipe 对象协作。
PipeWriter.GetMemory() 拿到一块可写内存,写完数据后调用 Advance(n) 告诉 Pipe "我写了 n 个字节",再调用 FlushAsync() 把数据提交给读取方,如果读取方处理太慢,FlushAsync 会根据背压配置异步等待;PipeReader.ReadAsync() 拿到一个 ReadResult,里面的 ReadOnlySequence<byte> 可能横跨多个内存块;ReadOnlySequence<byte> 里查找完整的消息边界(比如换行符),如果找到就处理,如果没找到就调用 AdvanceTo(consumed, examined),其中 examined 之前的数据"已检查但未消费",Pipe 会保留这部分等待更多数据到达;下图展示了 PipeWriter 与 PipeReader 之间基于内存段的协作流程:
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 选项可以做性能优化),常用于日志批处理、任务队列等场景。
如果流水线有多个处理阶段,且每个阶段的并行度、缓冲策略都不同,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 | 所在层 | 节点类型 | 典型场景 |
|---|---|---|---|
| ASP.NET Core 中间件管道 | Web 框架 | Middleware(RequestDelegate) | 鉴权、日志、路由、异常处理 |
| HttpClient 处理器链 | 出站 HTTP | DelegatingHandler | 重试、日志、熔断、签名 |
System.IO.Pipelines | 底层 I/O | PipeReader / PipeWriter | 自定义协议解析、高性能网络服务 |
Channel<T> | 并发编程 | Writer / Reader | 生产者-消费者、任务队列 |
| TPL Dataflow | 并发编程 | TransformBlock / ActionBlock 等 | 多阶段并行数据处理 |
IAsyncEnumerable + LINQ | 语言/运行时 | 迭代器 | 流式数据的惰性处理 |
DelegatingHandler;System.IO.Pipelines;Channel<T>;这些 Pipeline 并不是“二选一”的关系,而是不同层次的协作。一个典型的 Web 服务,通常会同时用到几种,但各自负责不同阶段。
HttpClient 的 DelegatingHandler 做统一重试、熔断、签名;Channel<T> 做缓冲队列,把“收到请求”与“后续异步处理”解耦;System.IO.Pipelines 处理字节流。可以把它理解成:
“请求进入 → 中间件做框架层横切 → HttpClient 做出站调用 → Channel<T> / TPL Dataflow 做异步工作流”。
System.IO.Pipelines 解析网络字节流;Channel<T> 做缓冲和背压控制;DelegatingHandler 调用上游服务;Channel<T> 接收实时事件;System.IO.Pipelines。DelegatingHandler 解决“出站请求的统一处理”;System.IO.Pipelines 解决“字节流/协议层的高性能处理”;Channel<T> / TPL Dataflow 解决“应用内异步、并行、复杂工作流”。.NET 里的各种 Pipeline 表面上互不相干,但本质都是责任链模式的变体:把一个大任务拆成若干个可以独立开发、独立测试、按需编排的小节点。理解了这套"进入 → 处理/传递 → 短路或返回"的通用流转逻辑,再看具体的中间件、消息处理器、PipeReader/PipeWriter,就只是"节点长什么样、数据是什么类型"的差异了。
暂无评论,快来抢沙发吧!