高性能按Key合并的生产者/消费者模式异步实现优化咨询
基于Key合并的异步工作队列实现优化
定义的接口
interface IConflateWorkByKey<TKey, TValue> { IAsyncEnumerable<Dictionary<TKey,TValue>> GetValues(); void Publish(TKey key, TValue value); }
该接口用于实现按Key合并的工作队列模式:生产者通过Publish方法提交任务项,单个消费者通过GetValues方法获取任务项,需满足以下要求:
- 按Key合并值:若生产者在消费者获取下一批任务(字典)前,先后调用
Publish(1, "hello")和Publish(1, "world"),消费者收到的下一个字典中Key为1的对应值应为"world",之前的更新(1, "hello")会被丢弃,不会被消费者感知。 Publish调用速度通常远快于任务消耗速度,因此Publish的调用需尽可能快,任务消耗的性能要求相对较低。- 多数实际场景中,消费者迭代
GetValues返回的迭代器时,新任务项已就绪,无需等待,需针对该场景做快速路径优化;同时实现也需支持无任务项时异步等待新任务的场景。 - 仅存在一个消费者(即
GetValues仅会被一个消费者调用/消费)。 Publish不会被并发调用(但可在不同线程中被顺序调用)。
当前实现及问题
class Conflate<TKey, TValue> : IConflateWorkByKey<TKey, TValue> { private Dictionary<TKey,TValue>? _buffered = null; private readonly object _lock = new(); public IAsyncEnumerable<Dictionary<TKey,TValue>> GetValues(CancellationToken ct) { while(!ct.IsCancellationRequested) { lock(_lock) { while(_buffered is null) Monitor.Wait(_lock); var result = _buffered; _buffered = null; yield return result; } } } public void Publish(TKey key, TValue value) { lock(_lock) { _buffered ??= new(); _buffered[key] = value; Monitor.Pulse(_lock); } } }
注:若某一实现更优,可将Publish方法改为返回ValueTask。
该实现原理上可行,但核心问题是GetValues并非真正异步,调用线程会被Monitor.Wait阻塞。
我也曾尝试使用Nito.AsyncEx中的AsyncMonitor实现该模式,但AsyncMonitor.Pulse的速度极慢,无法满足性能要求。
待验证的替代实现思路
class Conflate<TKey, TValue> : IConflateWorkByKey<TKey, TValue> { private Dictionary<TKey,TValue>? _buffered = new(); private readonly object _lock = new(); private TaskCompletionSource? _tcs = null; public IAsyncEnumerable<Dictionary<TKey,TValue>> GetValues(CancellationToken ct) { while(!ct.IsCancellationRequested) { Dictionary<TKey,TValue> result; while(true) { lock(_lock) { if(_buffered.Any()) { // 快速路径 - 结果已就绪,无需等待直接返回 result = _buffered; _buffered = new(); break; } _tcs = new(); } await _tcs.Task; } yield return result; } } public void Publish(TKey key, TValue value) { lock(_lock) { _buffered[key] = value; if(_tcs is not null) { _tcs.TrySetResult(); _tcs = null; // 快速路径:如果后续值未被取走,下一次Publish无需调用TrySetResult } } } }
需求目标
现寻求更巧妙的实现/模式,既要保证Publish调用的极致性能,又要支持GetValues内部真正的异步等待/通知机制。
内容的提问来源于stack exchange,提问作者Bogey
相关产品推荐
相关产品推荐

