System.Threading.Channels
使用 Channel、ChannelWriter 和 ChannelReader 构建生产者—消费者管道,以控制背压。
System.Threading.Channels 是 CoddyKit 上的免费 C# Academy 课时。 这是第 2 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 C# Academy 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 C# Academy 课程共包含 4 节课。
什么是通道?
System.Threading.Channels(在 .NET Core 3 中引入)提供高性能、线程安全的 producer-consumer 队列。与 BlockingCollection 不同,通道完全支持 async 且内存分配效率高,非常适合进程内管道。
创建通道
通道通过工厂创建。请选择 unbounded(无限制)或 bounded(容量有限并具有背压)。工厂会返回一个带有 Writer 和 Reader 的 Channel<T>。
using System.Threading.Channels;
// Unbounded: unlimited capacity, no backpressure
var unbounded = Channel.CreateUnbounded<string>();
// Bounded: max 100 items; writer waits when full
var bounded = Channel.CreateBounded<string>(100);
// Bounded with drop-oldest strategy:
var dropping = Channel.CreateBounded<string>(new BoundedChannelOptions(50)
{
FullMode = BoundedChannelFullMode.DropOldest
});向通道写入
producer 使用 WriteAsync 写入项目(bounded channel 已满时会等待),或使用 TryWrite(已满时立即返回 false)。使用 Writer.Complete() 发出完成信号。
var channel = Channel.CreateUnbounded<int>();
// Producer task
var producer = Task.Run(async () =>
{
for (int i = 0; i < 100; i++)
{
await channel.Writer.WriteAsync(i);
await Task.Delay(10);
}
channel.Writer.Complete(); // signal no more items
});从通道读取
consumer 可以使用 ReadAllAsync()(最简单的方式)读取,也可以使用 ReadAsync/TryRead 以获得更多控制权。当 writer 调用 Complete() 时,ReadAllAsync 会完成。
// Consumer task
var consumer = Task.Run(async () =>
{
await foreach (var item in channel.Reader.ReadAllAsync())
{
Console.WriteLine($"Processed: {item}");
// Naturally paced — waits for next item
}
Console.WriteLine("Channel completed");
});
await Task.WhenAll(producer, consumer);多个 producer
通道是线程安全的。多个 producer 可以并发写入,无需加锁。只有在所有 producer 完成后,才能调用 Writer.Complete()。
var channel = Channel.CreateUnbounded<WorkItem>();
var producers = Enumerable.Range(0, 4).Select(id =>
Task.Run(async () =>
{
for (int i = 0; i < 25; i++)
await channel.Writer.WriteAsync(new WorkItem(id, i));
}));
// Wait for all producers before completing the writer
await Task.WhenAll(producers);
channel.Writer.Complete();多个 consumer(扇出)
在同一 channel reader 上运行多个 consumer task,以并行处理数据。每个项目只会交付给一个 consumer(进行分区,而不是广播)。
var channel = Channel.CreateBounded<WorkItem>(100);
// 4 parallel consumers
var consumers = Enumerable.Range(0, 4).Select(id =>
Task.Run(async () =>
{
await foreach (var item in channel.Reader.ReadAllAsync())
{
await ProcessItemAsync(item);
Console.WriteLine($"Consumer {id} processed {item.Id}");
}
}));
await Task.WhenAll(consumers);管道模式
将多个 channel 串联成处理管道:每个阶段从一个 channel 读取数据,处理项目,然后写入下一个 channel。各阶段会并发运行,并自然地产生背压。
// Stage 1: raw data
var stage1 = Channel.CreateBounded<string>(50);
// Stage 2: parsed
var stage2 = Channel.CreateBounded<ParsedRecord>(50);
// Stage 3: enriched
var stage3 = Channel.CreateBounded<EnrichedRecord>(50);
var parse = ParseStageAsync(stage1.Reader, stage2.Writer);
var enrich = EnrichStageAsync(stage2.Reader, stage3.Writer);
var persist = PersistStageAsync(stage3.Reader);
await Task.WhenAll(parse, enrich, persist);使用有界通道实现背压
bounded channel 会自动施加背压:当 channel 已满时,WriteAsync 会暂停 producer,直到 consumer 取出项目。不需要手动编写限流代码。
// Bounded channel: max 10 items
var channel = Channel.CreateBounded<string>(10);
// Fast producer
var producer = Task.Run(async () =>
{
for (int i = 0; i < 1000; i++)
{
// WriteAsync waits when channel has 10 items
await channel.Writer.WriteAsync($"item-{i}");
// Producer is naturally slowed to consumer speed
}
channel.Writer.Complete();
});
// Slow consumer
var consumer = Task.Run(async () =>
{
await foreach (var item in channel.Reader.ReadAllAsync())
{
await Task.Delay(50); // simulate slow processing
Console.WriteLine(item);
}
});通道中的错误处理
将异常传递给 Writer.Complete(exception),即可将错误传播给所有正在等待的读取方。读取方会在下一次读取 channel 时看到该异常。
var channel = Channel.CreateUnbounded<int>();
var producer = Task.Run(async () =>
{
try
{
for (int i = 0; i < 100; i++)
{
if (i == 50) throw new Exception("Producer failed at 50");
await channel.Writer.WriteAsync(i);
}
channel.Writer.Complete();
}
catch (Exception ex)
{
channel.Writer.Complete(ex); // propagate to reader
}
});
try
{
await foreach (var item in channel.Reader.ReadAllAsync())
Console.WriteLine(item);
}
catch (Exception ex)
{
Console.Error.WriteLine($"Channel error: {ex.Message}");
}实际案例:后台作业队列
使用托管服务和 channel 构建后台作业队列:HTTP 请求将作业加入队列,后台工作器逐个处理作业。
public class JobQueue
{
private readonly Channel<Func<CancellationToken, Task>> _queue
= Channel.CreateBounded<Func<CancellationToken, Task>>(100);
public ChannelWriter<Func<CancellationToken, Task>> Writer => _queue.Writer;
public ChannelReader<Func<CancellationToken, Task>> Reader => _queue.Reader;
}
public class JobProcessor : BackgroundService
{
private readonly JobQueue _queue;
public JobProcessor(JobQueue q) => _queue = q;
protected override async Task ExecuteAsync(CancellationToken ct)
{
await foreach (var job in _queue.Reader.ReadAllAsync(ct))
await job(ct);
}
}快速检查
当 producer 在已满的 bounded channel 上调用 WriteAsync 时会发生什么?
回顾:System.Threading.Channels
要点:
- channel 提供高性能、原生支持 async 的 producer-consumer 队列
- unbounded:无限制;bounded:容量有限,并支持背压或丢弃策略
- 多个 producer 和 consumer 无需加锁即可安全运行
- ReadAllAsync() + await foreach 是最简洁的消费模式
- 将 channel 串联成管道,以进行基于阶段的并发处理
- 调用 Writer.Complete()(或 Complete(exception))发出流结束信号
常见问题解答
「System.Threading.Channels」课时是免费的吗?
是的 — 「System.Threading.Channels」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 C# Academy 课程的其余内容,请升级到 CoddyKit PRO。 C# Academy 课程共包含 4 节课。
「System.Threading.Channels」这节课中我会学到什么?
使用 Channel、ChannelWriter 和 ChannelReader 构建生产者—消费者管道,以控制背压。 你通过在浏览器中直接运行的动手代码来练习 C# Academy,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 C# Academy 需要有经验吗?
无需任何先前经验。CoddyKit 上的 C# Academy 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 2 节课,共 4 节。
「System.Threading.Channels」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 C# Academy 课中编写并运行代码吗?
能。每节 C# Academy 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。
此课程中的所有课时
- IAsyncEnumerable 与 await foreach
- System.Threading.Channels
- ValueTask 与避免分配
- ConfigureAwait 与同步上下文