0Pricing
C# Academy · 课时

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 反馈 — 无需本地设置。

此课程中的所有课时

  1. IAsyncEnumerable 与 await foreach
  2. System.Threading.Channels
  3. ValueTask 与避免分配
  4. ConfigureAwait 与同步上下文
← 返回 C# Academy