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

C# .NET Channel多消费者不丢消息 各消费者接收全量消息方案

问题结论

你当前的代码完全无法实现「所有消费者接收全量消息、不丢消息」的需求,存在核心逻辑错误。

现有代码的问题
  • 核心模型错误:System.Threading.Channels 是单播队列模型,无论是否设置SingleReader = false,一条消息写入后只会被一个消费者读取到,多个消费者是竞争取消息的关系,根本不会给每个消费者推送消息副本。运行现有代码会发现两个Reader交替拿到数字,没有任何一个Reader能拿到全量消息。
  • 写入逻辑错误:你设置队列容量仅为1,且使用TryWrite写入。TryWrite在队列满时会直接返回false不会等待,只要消费者没有及时消费,消息会直接丢弃,完全不符合「不丢消息」的要求。另外SingleReader = false只是给运行时的性能优化提示,不会改变队列的单播分发逻辑。
正确实现方案

你的场景属于典型的广播推送场景,正确的实现逻辑不是让多个消费者共享同一个Channel,而是给每个WebSocket连接分配独立的专属消费队列:

  1. 维护一个线程安全的集合,存储当前所有在线WebSocket连接对应的专属Channel
  2. 从Redis Pub/Sub收到消息后,遍历所有连接的专属Channel,给每个Channel写入一份消息副本
  3. 每个WebSocket连接仅消费自己专属Channel的消息,发送给客户端,不会和其他连接竞争消息
  4. 连接断开时及时从集合中移除对应的Channel,完成资源释放

参考实现代码

// 并发字典存储所有连接的专属Channel,Key可替换为你自己的连接唯一标识
var connectionChannels = new ConcurrentDictionary<string, Channel<int>>();
var cts = new CancellationTokenSource();

// 模拟Redis Pub/Sub消息接收+广播逻辑
var broadcastTask = Task.Run(async () =>
{
    int msgSeq = 0;
    while (!cts.IsCancellationRequested)
    {
        var currentMsg = msgSeq++;
        // 给每个在线连接的队列都写入一份消息
        foreach (var (connId, channel) in connectionChannels)
        {
            try
            {
                // 写入设置100ms超时,避免慢客户端阻塞整体广播流程
                await channel.Writer.WriteAsync(currentMsg, cts.Token)
                    .AsTask()
                    .WaitAsync(TimeSpan.FromMilliseconds(100), cts.Token);
            }
            catch (TimeoutException)
            {
                // 写入超时说明客户端接收过慢,主动移除队列、断开连接
                if (connectionChannels.TryRemove(connId, out var slowChannel))
                {
                    slowChannel.Writer.Complete();
                    // 此处补充你的WebSocket关闭逻辑
                }
            }
        }
        await Task.Delay(TimeSpan.FromMilliseconds(250), cts.Token);
    }

    // 广播停止时关闭所有连接的队列
    foreach (var channel in connectionChannels.Values)
    {
        channel.Writer.Complete();
    }
});

// 单连接处理逻辑,每个WebSocket连接接入时调用一次
async Task HandleConnectionAsync(string connId, CancellationToken ct)
{
    // 为当前连接创建独立的有界队列,容量根据业务吞吐量、可接受延迟调整
    var ownChannel = Channel.CreateBounded<int>(new BoundedChannelOptions(100)
    {
        FullMode = BoundedChannelFullMode.Wait,
        SingleReader = true, // 仅当前连接读,开启性能优化
        SingleWriter = true  // 仅广播任务写,开启性能优化
    });
    connectionChannels.TryAdd(connId, ownChannel);

    try
    {
        // 仅消费自己的队列,发送给客户端
        await foreach (var msg in ownChannel.Reader.ReadAllAsync(ct))
        {
            Console.WriteLine($"{connId} 收到消息: {msg}");
            // 此处补充实际的WebSocket消息发送逻辑
        }
    }
    finally
    {
        // 连接断开时清理资源
        connectionChannels.TryRemove(connId, out _);
        ownChannel.Writer.Complete();
    }
}

// 启动两个模拟连接
var conn1Task = HandleConnectionAsync("ReaderOne", cts.Token);
var conn2Task = HandleConnectionAsync("ReaderTwo", cts.Token);

// 5秒后停止测试
cts.CancelAfter(TimeSpan.FromSeconds(5));
await Task.WhenAll(broadcastTask, conn1Task, conn2Task);

生产环境注意事项

  • 每个连接的队列必须设置合理容量,不要用无界队列,避免慢客户端导致内存溢出
  • 给队列写入设置合理超时,超时的慢客户端直接断开,避免阻塞整体广播流程
  • 不要在广播遍历过程中做耗时操作,所有IO逻辑(比如实际WebSocket发送)放到对应连接的独立消费任务里执行

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 05:27:31