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

如何使用TaskCompletionSource<T>实现无锁异步等待数据采样?

问题

我有一个处理无限数据流的类,外部组件需不定期异步采样处理后的数据。因下一条数据可能需等待较久,希望检查组件通过await而非阻塞等待,多个检查器可获取同一数据实例。
数据项体积庞大,为保证内存效率,处理后会立即释放,且不想做任何延迟释放的操作(如存入通道等)。
目前我使用TaskCompletionSource结合锁实现了该功能,代码如下:

public class DataProcessor
{
  private object _sampleDataLock = new object();
  private TaskCompletionSource<Data> _sampleData;
  private Thread _processThread;

  public DataProcessor()
  {
     _processThread = new(ProcessDataLoop);
     _processThread.Start();
  }

  void ProcessDataLoop()
  {
      while (true)
      {
         var processedData = ProcessNextData();
         
         //This code in the lock is what I wonder about
         lock (_sampleDataLock)
         {
            _sampleData?.SetResult(processedData.Clone());
            _sampleData = null;
         }

        processedData.Dispose();
      }
  }

  public Task<Data> SampleData(TimeSpan timeout)
  {
    // This is the code I'm wondering how to write without a lock
    lock (_sampleDataLock)
    {
      _sampleData ??= new TaskCompletionSource<Data>();
      return _sampleData.Task.WithTimeout(timeout);
    }
  }
}

public class DataInspector
{
    private DataProcessor _processor;

    public async Task InspectData()
    {
       //wait up to 30s for an item
       var data = await _processor.SampleData(TimeSpan.FromSeconds(30));
       // look at data
    }
}

该实现可正常运行,但我希望改为无锁实现,请问是否有可行方案?

无锁实现方案

可以用Interlocked类的原子操作替代锁,核心是利用Interlocked.CompareExchange和Interlocked.Exchange实现线程安全的TaskCompletionSource<Data>访问与替换,同时保留原有的内存效率。

核心思路

  1. 用原子操作替代锁逻辑,避免线程阻塞
  2. 保证数据生成线程与采样线程对_sampleData的访问无竞态条件
  3. 维持数据处理完成后立即释放的逻辑,不引入额外内存开销

无锁实现代码

public class DataProcessor
{
    // volatile关键字保证多线程下变量可见性,避免缓存旧值
    private volatile TaskCompletionSource<Data> _sampleData;
    private Thread _processThread;

    public DataProcessor()
    {
        _processThread = new(ProcessDataLoop);
        _processThread.Start();
    }

    void ProcessDataLoop()
    {
        while (true)
        {
            var processedData = ProcessNextData();
            
            // 原子替换:获取当前TCS并将_sampleData设为null
            var currentTcs = Interlocked.Exchange(ref _sampleData, null);
            currentTcs?.SetResult(processedData.Clone());

            processedData.Dispose();
        }
    }

    public Task<Data> SampleData(TimeSpan timeout)
    {
        while (true)
        {
            var currentTcs = _sampleData;
            // 如果已有等待中的TCS,直接返回其任务
            if (currentTcs != null)
            {
                return currentTcs.Task.WithTimeout(timeout);
            }

            // 创建新TCS,尝试原子替换空值
            var newTcs = new TaskCompletionSource<Data>();
            var originalTcs = Interlocked.CompareExchange(ref _sampleData, newTcs, null);
            // 替换成功则返回新TCS任务,失败则说明其他线程已创建TCS,循环重试
            if (originalTcs == null)
            {
                return newTcs.Task.WithTimeout(timeout);
            }
        }
    }
}

// 原DataInspector类无需修改
public class DataInspector
{
    private DataProcessor _processor;

    public async Task InspectData()
    {
        var data = await _processor.SampleData(TimeSpan.FromSeconds(30));
        // 处理数据
    }
}

关键细节说明

  • volatile关键字:确保_sampleData的更新能被所有线程即时感知,避免线程因缓存旧值导致逻辑错误
  • Interlocked.Exchange:数据生成线程中原子性获取并清空_sampleData,保证所有等待的采样线程都能收到数据信号
  • Interlocked.CompareExchange:采样线程中仅当_sampleData为空时才创建新TCS,避免多个线程重复创建,保证同一时刻只有一个TCS等待新数据
  • 循环重试逻辑:若创建新TCS时发现已有其他线程创建,则重新获取已存在的TCS,规避竞态条件

这种实现完全移除了锁,同时保留了原有的内存高效性——数据处理完成后立即释放,无额外缓存,多个采样线程也能获取到同一数据实例的克隆。

内容的提问来源于stack exchange,提问作者John V

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 21:30:23