如何使用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>访问与替换,同时保留原有的内存效率。
核心思路
- 用原子操作替代锁逻辑,避免线程阻塞
- 保证数据生成线程与采样线程对
_sampleData的访问无竞态条件 - 维持数据处理完成后立即释放的逻辑,不引入额外内存开销
无锁实现代码
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
相关产品推荐
相关产品推荐

