You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

高性能按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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.07 13:00:26