在C#服务端开发中,经常需要把海量数据按阶段依次处理,例如从消息队列拉取原始报文,经过反序列化、业务校验、聚合计算,最后写入数据库。如果采用单一线程顺序执行,CPU和IO资源都无法被充分利用。编写可扩展的并发数据处理管道,本质是把处理流程拆成多个阶段,每个阶段可以并行运行多个工作实例,阶段之间通过线程安全的数据结构传递数据,从而让整体吞吐随机器核数线性增长。

为什么选择Channel而非传统队列
早期不少项目用ConcurrentQueue配合Task自己写调度,但这种方式要手动处理完成通知、背压和异常传播。一旦某个阶段消费太慢,队列会无限堆积内存。C#自.NET Core 3.0起提供的System.Threading.Channels库,专门解决生产者消费者通信问题。它内置背压:当通道容量满时,生产者写入会自动等待,不需要业务代码控制限流。
Channel分为有界和无界两种。有界通道通过CreateBounded设定最大元素数,天然实现背压;无界通道用CreateUnbounded,适合内存充裕且下游处理极快的场景。在数据处理管道中,几乎总是推荐使用有界通道,避免突发流量打挂进程。另外,Channel支持多生产者多消费者,每个阶段开启若干个后台任务读取同一个通道,就能轻松横向扩展。
管道的基础结构定义
一个典型管道包含源读取、中间处理、落库三个环节。我们用泛型通道串联,每个环节的输入输出类型可以不同。下面代码展示如何创建有界通道并启动多消费者:
using System.Threading.Channels;
// 创建有界通道,容量为1000,满时写入方等待
var rawChannel = Channel.CreateBounded<string>(new BoundedChannelOptions(1000)
{
FullMode = BoundedChannelFullMode.Wait
});
// 启动3个解析任务
for (int i = 0; i < 3; i++)
{
_ = Task.Run(async () =>
{
await foreach (var item in rawChannel.Reader.ReadAllAsync())
{
// 解析逻辑
}
});
}
上面的ReadAllAsync会在通道关闭且数据读完时自动结束循环,不需要额外退出标志。这种写法比自己维护CancellationToken加轮询要简洁得多。阶段之间传递的模型建议用不可变记录(record),避免多消费者并发修改同一对象引发隐蔽bug。
多阶段串联与异常隔离
真实管道往往不止两个阶段。我们可以把每个阶段封装成方法,接收上游ChannelReader,返回下游ChannelWriter。这样主程序像搭积木一样组合。关键是异常隔离:某个数据导致处理失败不应让整个管道崩溃,应在阶段内捕获并记入错误通道。
async Task RunStageAsync<TIn, TOut>(
ChannelReader<TIn> input,
ChannelWriter<TOut> output,
Func<TIn, Task<TOut>> processor)
{
try
{
await foreach (var data in input.ReadAllAsync())
{
try
{
var result = await processor(data);
await output.WriteAsync(result);
}
catch (Exception ex)
{
// 单条数据异常,记录后继续
Console.WriteLine($"处理失败: {ex.Message}");
}
}
}
finally
{
output.TryComplete();
}
}
上述代码中,外层try保证即使上游异常中断,本阶段也会调用TryComplete通知下游。内层try保证单条脏数据不影响同阶段其他任务。这种双层结构在扩容到几十个消费者时依然稳健。
动态扩展与完成信号
可扩展性不仅指代码能加线程,还包括运行时按需调整。由于Channel的消费者就是普通Task,我们可以在监控指标显示积压时,动态启动更多消费任务;当积压清空,多余任务在通道完成后自然退出。完成信号由源端调用Writer.Complete触发,随后像多米诺骨牌一样逐阶段关闭。
| 阶段 | 建议并发数 | 扩容依据 |
|---|---|---|
| 读取源 | 1-2 | 受限于源头吞吐 |
| CPU计算 | 核数×1 | CPU使用率 |
| IO落库 | 核数×2 | 数据库响应时间 |
表中给出常见阶段的并发参考。IO密集型阶段因为大部分时间在等待网络,可以开更多任务提升利用率。但要注意数据库连接池上限,避免管道扩容把底层资源压垮。配合SemaphoreSlim限制对外部系统的并发调用,是生产环境常用做法。
完整示例骨架
下面把前面内容拼成一个最小可运行骨架,展示从字符串读取到落库的三阶段管道。你可以直接复制到项目里改成自己的业务逻辑。
using System.Threading.Channels;
var source = Channel.CreateBounded<string>(100);
var parsed = Channel.CreateBounded<int>(100);
var saved = Channel.CreateBounded<int>(100);
// 阶段1:模拟数据源
_ = Task.Run(async () =>
{
for (int i = 0; i < 1000; i++)
await source.Writer.WriteAsync(i.ToString());
source.Writer.Complete();
});
// 阶段2:解析
_ = RunStageAsync(source.Reader, parsed.Writer, async s =>
{
await Task.Yield();
return int.Parse(s);
});
// 阶段3:保存
_ = RunStageAsync(parsed.Reader, saved.Writer, async v =>
{
await Task.Delay(1);
return v;
});
// 收尾
_ = Task.Run(async () =>
{
await foreach (var _ in saved.Reader.ReadAllAsync()) { }
Console.WriteLine("管道处理完毕");
});
async Task RunStageAsync<TIn, TOut>(ChannelReader<TIn> input,
ChannelWriter<TOut> output, Func<TIn, Task<TOut>> processor)
{
try
{
await foreach (var data in input.ReadAllAsync())
{
try { await output.WriteAsync(await processor(data)); }
catch (Exception e) { Console.WriteLine(e.Message); }
}
}
finally { output.TryComplete(); }
}
这个骨架里每个阶段都通过RunStageAsync启动,源端完成后信号自动向后传递。如果想扩展解析阶段,只需在调用RunStageAsync时多开几个任务读同一个source.Reader即可,不需要改动其他阶段代码。这就是可扩展并发管道在C#中的核心写法。