如何在同一线程高效读取两个不同的System.Threading.Channels?
优先级事件流处理方案咨询
场景说明
- 存在两种不同事件流
- 流1:基于拉取的高速事件流(如金融工具价格)
- 流2:低速事件流(如基于价格的操作指令)
- 需在同一工作线程/任务中处理两种流
- 实现需尽可能高效
- 流2必须优先于流1处理
现有尝试的问题
将所有事件推入System.Threading.Channel<AnyOf<Event1, Event2>>的方案能满足多数要求,但无法实现流2的优先级处理。
尝试用独立通道分别推送事件并同时读取,示例代码如下:
class Example { private readonly Channel<Event1> _channel1 = Channel.CreateBounded<Event1>(....); private readonly Channel<Event2> _channel2 = Channel.CreateBounded<Event2>(....); public async Task Worker(CancellationToken ct) { while(!ct.IsCancellationRequested) { var valueTask1 = _channel1.Reader.WaitToReadAsync(ct); var valueTask2 = _channel2.Reader.WaitToReadAsync(ct); // 问题来了!.NET没有内置的ValueTask.WhenAny()方法 if(await ValueTaskExtensions.WhenAny(valueTask1, valueTask2) == 0) { // 处理流1事件 var element = await _channel1.Reader.ReadAsync(ct); // 处理逻辑 } else { // 处理流2事件 var element = await _channel2.Reader.ReadAsync(ct); // 处理逻辑 } } } }
但.NET没有内置的ValueTask.WhenAny(),如果转成Task会因为流1的高频事件产生大量内存分配开销。虽然社区有第三方实现的WhenAny,但未实际使用过,无法评估其质量和风险。
现在想咨询:有没有更巧妙的方法,能在同一读取线程中读取两个通道,且全程不用脱离ValueTask的使用?
解决方案
方案一:绝对优先轮询(最简洁高效)
核心逻辑是先清空流2的所有待处理事件,再处理流1的单个事件,完全不需要WhenAny,全程基于ValueTask操作,保证流2的绝对优先级。
class PriorityEventProcessor { private readonly Channel<Event1> _highSpeedChannel = Channel.CreateBounded<Event1>(new BoundedChannelOptions(1000) { SingleReader = true }); private readonly Channel<Event2> _priorityChannel = Channel.CreateBounded<Event2>(new BoundedChannelOptions(100) { SingleReader = true }); public async Task WorkerLoop(CancellationToken ct) { while (!ct.IsCancellationRequested) { // 优先处理所有流2的待处理事件 while (await _priorityChannel.Reader.WaitToReadAsync(ct)) { if (_priorityChannel.Reader.TryRead(out var priorityEvent)) { ProcessPriorityEvent(priorityEvent); } } // 流2无数据时,处理一个流1的事件 if (await _highSpeedChannel.Reader.WaitToReadAsync(ct)) { if (_highSpeedChannel.Reader.TryRead(out var highSpeedEvent)) { ProcessHighSpeedEvent(highSpeedEvent); } } } } private void ProcessPriorityEvent(Event2 evt) { // 流2事件处理逻辑 } private void ProcessHighSpeedEvent(Event1 evt) { // 流1事件处理逻辑 } }
这种方式的优势是:
- 完全避免
WhenAny带来的内存分配问题 - 流2优先级绝对保障,只要有未处理事件就会持续处理,不会被流1抢占
- 配置
SingleReader进一步提升通道读取效率
方案二:双通道异步等待优化
如果需要在流2无数据时同时等待两个通道的新事件,又不想产生过多Task分配,可以用以下方式:
public async Task WorkerLoop(CancellationToken ct) { while (!ct.IsCancellationRequested) { var priorityWait = _priorityChannel.Reader.WaitToReadAsync(ct); // 先检查流2是否有数据就绪 if (priorityWait.IsCompletedSuccessfully) { if (_priorityChannel.Reader.TryRead(out var evt)) { ProcessPriorityEvent(evt); continue; } } // 流2无就绪数据,同时等待两个通道 var highSpeedWait = _highSpeedChannel.Reader.WaitToReadAsync(ct); // 仅在必要时转为Task,减少高频场景下的分配 var completedTask = await Task.WhenAny(priorityWait.AsTask(), highSpeedWait.AsTask()); if (completedTask == priorityWait.AsTask()) { if (await priorityWait && _priorityChannel.Reader.TryRead(out var evt)) { ProcessPriorityEvent(evt); } } else { if (await highSpeedWait && _highSpeedChannel.Reader.TryRead(out var evt)) { ProcessHighSpeedEvent(evt); } } } }
这个方案在流2频繁有数据时,大部分逻辑用ValueTask完成,只有双通道等待时才会短暂使用Task,平衡了优先级和性能。
内容的提问来源于stack exchange,提问作者Bogey
相关产品推荐
相关产品推荐

