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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 10:17:58