如何对缓冲的Observable令牌按优先级排序并高效处理?
针对你的问题,我们可以结合TPL Dataflow和一个自定义的线程安全优先级缓冲区来实现优雅的排序缓冲+异步处理,既避免轮询,又解决多消费者的竞态问题。你的Token类已经实现了IComparable<Token>,正好可以利用这个排序逻辑来维护缓冲区的顺序。
步骤1:实现线程安全的优先级缓冲区
这个缓冲区用SortedSet<Token>维护有序的令牌,通过TaskCompletionSource实现异步等待(避免轮询),并用锁保证线程安全:
public class PriorityBuffer { private readonly SortedSet<Token> _sortedTokens = new SortedSet<Token>(); private readonly object _lock = new object(); private TaskCompletionSource<bool> _tcs = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously); private bool _isCompleted; // 添加令牌到缓冲区 public void Add(Token token) { if (_isCompleted) throw new InvalidOperationException("缓冲区已完成,无法添加新令牌"); lock (_lock) { _sortedTokens.Add(token); _tcs.TrySetResult(true); // 通知等待的任务有新元素 } } // 标记缓冲区完成(不再接收新元素) public void Complete() { lock (_lock) { _isCompleted = true; _tcs.TrySetResult(true); // 唤醒所有等待的任务 } } // 异步获取下一个最高优先级的令牌(无元素时等待) public async Task<Token?> TakeAsync(CancellationToken cancellationToken = default) { while (true) { cancellationToken.ThrowIfCancellationRequested(); lock (_lock) { // 如果有元素,取出优先级最高的(SortedSet的First()对应你的CompareTo逻辑的最高优先级) if (_sortedTokens.Count > 0) { var token = _sortedTokens.First(); _sortedTokens.Remove(token); return token; } // 缓冲区已完成且无元素,返回null结束 if (_isCompleted) return null; // 重置等待源,准备等待新元素 _tcs = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously); } // 异步等待新元素或缓冲区完成 await _tcs.Task.WaitAsync(cancellationToken); } } }
步骤2:整合生产者、缓冲区和处理器
修改你的生产者代码,将生成的令牌送入优先级缓冲区,再用TPL Dataflow的ActionBlock处理令牌(控制并发度):
// 1. 创建处理器:模拟慢处理,可设置并发度 var processor = new ActionBlock<Token>(async token => { // 模拟耗时处理(替换成你的实际逻辑) await Task.Delay(500); Console.ForegroundColor = ConsoleColor.Green; Console.WriteLine($"{DateTime.Now:mm:ss.fff}:已处理 {token}"); }, new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 1, // 可根据需要调整并发数,若处理逻辑线程安全可设为更高值 BoundedCapacity = 10 // 设置处理器的缓冲区大小,控制背压 }); // 2. 创建优先级缓冲区 var priorityBuffer = new PriorityBuffer(); // 3. 订阅令牌流,将令牌加入缓冲区 var subscription = tokens.Subscribe( token => priorityBuffer.Add(token), () => priorityBuffer.Complete() // 令牌生成完毕,标记缓冲区完成 ); // 4. 启动异步消费任务:从缓冲区取令牌并送入处理器 _ = Task.Run(async () => { try { while (true) { var token = await priorityBuffer.TakeAsync(); if (token == null) break; // 缓冲区已完成且无元素,退出循环 // 发送令牌到处理器(若处理器缓冲区满则自动等待,处理背压) await processor.SendAsync(token); } } catch (OperationCanceledException) { // 处理取消逻辑(可选) } finally { processor.Complete(); // 通知处理器不再接收新任务 subscription.Dispose(); // 释放订阅资源 } }); // 5. 等待所有处理完成 await processor.Completion; Console.ResetColor(); Console.WriteLine("所有令牌处理完成");
方案优势
- 无轮询:通过
TaskCompletionSource实现异步等待,有新元素时才唤醒消费逻辑,避免空循环浪费资源。 - 线程安全:用锁保护
SortedSet的操作,多消费者场景下也不会出现竞态条件(若需要多消费任务,直接启动多个即可,TakeAsync是线程安全的)。 - 背压控制:TPL Dataflow的
ActionBlock支持BoundedCapacity,当处理器忙时会阻塞SendAsync,间接控制生产者的速度(若需要更严格的背压,可调整缓冲区参数)。 - 优先级排序:利用你已实现的
IComparable<Token>逻辑,自动按High>Medium>Low排序,同优先级按对应子类的Id排序。
内容的提问来源于stack exchange,提问作者Aaron
相关产品推荐
相关产品推荐

