C# Channel多生产者消费者场景如何实现生产与消费并行运行
问题解决方案
你遇到的问题本质是LINQ查询的惰性枚举特性导致的:你调用Enumerable.Select返回的IEnumerable<Task>是惰性计算的,只有当你第一次枚举它的时候才会实际执行里面的Task.Run逻辑创建并启动任务。
你原来的代码里,直到执行await Task.WhenAll(producers)的时候才会第一次枚举producers启动生产者任务,而消费者任务更是要等到执行await Task.WhenAll(consumers)的时候才会被枚举启动,自然就出现了生产完才消费的情况。
调整方案
只需要两处改动即可实现生产消费并行:
- 把producers和consumers的LINQ查询结果立刻转成数组/列表,强制触发枚举,启动所有任务
- 把
ContinueWith换成直接await后调用Complete,更符合异步编程最佳实践,避免异常捕获、线程调度的潜在问题
修改后完整代码
var channel = Channel.CreateUnbounded<string>(); // 加.ToArray() 立刻启动所有消费者任务 var consumers = Enumerable .Range(1, 5) .Select(consumerNumber => Task.Run(async () => { var rnd = new Random(); while (await channel.Reader.WaitToReadAsync()) { if (channel.Reader.TryRead(out var item)) { Console.WriteLine($"Consuming {item} on consumer {consumerNumber}"); } } })) .ToArray(); // 加.ToArray() 立刻启动所有生产者任务 var producers = Enumerable .Range(1, 5) .Select(producerNumber => Task.Run(async () => { var rnd = new Random(); for (var i = 0; i < 10; i++) { var t = $"Message {i} from producer {producerNumber}"; Console.WriteLine($"Producing {t} on producer {producerNumber}"); await channel.Writer.WriteAsync(t); await Task.Delay(TimeSpan.FromSeconds(rnd.Next(3))); } })) .ToArray(); // 等所有生产者生产完毕后关闭通道 await Task.WhenAll(producers); channel.Writer.Complete(); // 等所有消费者把剩余消息消费完 await Task.WhenAll(consumers);
调整后效果
所有生产者和消费者任务会在你调用.ToArray()的时候同时启动,生产者每写入一条消息,空闲的消费者就会立刻读取消费,完全并行运行。
内容的提问来源于stack exchange,提问作者Houlahan
相关产品推荐
相关产品推荐

