如何在C#异步编程中实现多读者-单写者的读写任务同步(唤醒所有等待读者场景)
如何在C#异步编程中实现多读者-单写者的读写任务同步(唤醒所有等待读者场景)
老兄,我太懂你这种感受了——搜了一堆异步教程全是“煮咖啡”“做早餐”这种入门例子,真要搞这种多等待者唤醒的生产场景,反而找不到精准的方案,尤其是在Kestrel这种Web环境下,还要避免线程池被耗尽的坑。
你要的其实就是异步版的Monitor.PulseAll,核心是用TaskCompletionSource<T>来实现无阻塞的等待,同时维护一个等待列表,每次有新数据时唤醒所有等待的任务。下面直接上可落地的代码和解释,完全不需要额外库,适配你的Kestrel场景:
核心实现:DataBroadcaster类
这个类负责管理等待的读者任务,以及处理写者的新数据发布:
public class DataBroadcaster<T> { // 保护共享状态的锁 private readonly object _lockObj = new object(); // 存储所有等待新数据的任务源 private readonly List<TaskCompletionSource<bool>> _waiters = new List<TaskCompletionSource<bool>>(); // 最新的待发布数据 private T _latestData; // 标记是否有未被取走的新数据 private bool _hasNewData; /// <summary> /// 读者调用:等待新数据或超时返回 /// </summary> /// <param name="timeout">等待超时时间</param> /// <returns>新数据,超时则返回默认值</returns> public async Task<T?> WaitForNewDataAsync(TimeSpan timeout) { lock (_lockObj) { // 如果已有未被取走的新数据,直接返回并重置标记 if (_hasNewData) { var data = _latestData; _hasNewData = false; return data; } // 没有新数据,创建任务源加入等待列表 var tcs = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously); _waiters.Add(tcs); // 这里不能在lock内await,所以用Task.WhenAny处理等待/超时逻辑 return await Task.WhenAny(tcs.Task, Task.Delay(timeout)).ContinueWith(completedTask => { lock (_lockObj) { // 不管是完成还是超时,都从等待列表移除当前任务源,避免内存泄漏 _waiters.Remove(tcs); } // 检查是否是等待任务完成(而非超时) if (completedTask.Result == tcs.Task && tcs.Task.IsCompletedSuccessfully) { lock (_lockObj) { var data = _latestData; _hasNewData = false; return data; } } // 超时或任务失败,返回默认值 return default(T); }, TaskContinuationOptions.ExecuteSynchronously); } } /// <summary> /// 写者调用:发布新数据并唤醒所有等待的读者 /// </summary> /// <param name="data">要发布的新数据</param> public void PublishNewData(T data) { lock (_lockObj) { // 更新最新数据并标记为有新数据 _latestData = data; _hasNewData = true; // 唤醒所有等待的读者任务 foreach (var waiter in _waiters) { waiter.TrySetResult(true); } // 清空等待列表,因为所有等待者已被唤醒 _waiters.Clear(); } } }
关键细节解释
- 线程安全保障:所有对共享状态(等待列表、最新数据、新数据标记)的操作都用
lock包裹,确保在多线程环境下不会出现竞态条件。 - 无阻塞等待:用
TaskCompletionSource实现异步等待,读者任务在await时会释放线程回到线程池,不会像传统的Monitor.Wait那样占用线程,完美适配Kestrel的高并发场景。 - 超时处理:通过
Task.WhenAny结合Task.Delay实现超时逻辑,超时后会自动从等待列表移除对应的任务源,避免内存泄漏。 - 异步延续优化:创建
TaskCompletionSource时指定TaskCreationOptions.RunContinuationsAsynchronously,确保唤醒读者的延续任务在异步线程执行,不会阻塞写者所在的线程(比如Kestrel的IO线程)。
在Kestrel中的使用示例
你可以直接在API控制器中注入或实例化DataBroadcaster,实现等待型接口和数据发布接口:
[ApiController] [Route("api/data")] public class DataController : ControllerBase { // 全局单例的DataBroadcaster,也可以用依赖注入注册为单例 private static readonly DataBroadcaster<string> _dataBroadcaster = new DataBroadcaster<string>(); /// <summary> /// 写者接口:发布新数据 /// </summary> [HttpPost("publish")] public IActionResult PublishNewData([FromBody] string data) { _dataBroadcaster.PublishNewData(data); return Ok("Data published successfully"); } /// <summary> /// 读者接口:等待新数据,超时10秒 /// </summary> [HttpGet("wait")] public async Task<IActionResult> WaitForNewData() { var newData = await _dataBroadcaster.WaitForNewDataAsync(TimeSpan.FromSeconds(10)); if (newData == null) { return Ok("No new data received within timeout"); } return Ok($"Received new data: {newData}"); } }
扩展与注意事项
- 支持多数据项:如果需要保存多个数据项而非仅最新一条,可以把
_latestData改成Queue<T>,读者进来先检查队列是否有数据,有就出队返回,没有再等待;写者入队后唤醒所有等待者。 - 内存泄漏防范:务必确保所有
TaskCompletionSource都被处理(完成或超时后从列表移除),避免长期占用内存。 - 异常处理:如果需要处理任务失败场景,可以在
WaitForNewDataAsync中加入TrySetException的逻辑,根据你的业务需求调整。
内容来源于stack exchange
相关产品推荐
相关产品推荐

