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

如何使用System.Reactive实现多消费者无遗漏接收全量消息

回答

你写的示例核心思路是对的,但直接套到生产环境的Redis+WebSocket场景会踩几个致命坑,没法满足“所有已连接WebSocket都能收到全部后续消息、不遗漏、不互相影响”的要求。

现有示例代码的问题

你用的Publish()返回的是普通可连接可观察序列,行为逻辑是:手动调用Connect()之后序列才开始生产消息,只有Connect()执行前完成订阅的消费者能拿到全量消息。

  • 只要Connect()执行完成后新接入的WebSocket,会直接漏掉订阅前已经发过的所有消息
  • 没有异常隔离:只要某一个WebSocket推送时抛异常(比如连接中途断开没及时清理),整个上游序列会直接终止,剩下所有消费者都收不到后续消息
  • 没有自动生命周期管理:WebSocket断开后如果不手动释放订阅句柄,会一直占用内存造成泄漏,也不会自动断开闲置的Redis订阅浪费连接资源。

适配业务场景的正确实现

你的场景核心要求很明确:

  • Redis Pub/Sub 整个进程只需要建1个订阅,不要重复建连浪费资源
  • 任意时刻新接入的WebSocket,从接入时刻开始必须收到所有后续到达的Redis消息,不能丢
  • 单个WebSocket推送失败不能牵连其他连接
  • WebSocket断开后自动取消订阅,清理资源

直接参考下面的实现逻辑即可:

// 这部分全局只初始化一次,作为整个进程的Redis消息广播源
var redisMessageStream = Observable.Create<YourBusinessMessage>(async (observer, cancelToken) =>
{
    // 替换成你实际用的Redis客户端订阅逻辑,比如StackExchange.Redis的订阅实现
    var redisSub = await redisClient.GetSubscriber().SubscribeAsync("your_message_channel");
    redisSub.OnMessage(rawMsg =>
    {
        var msg = DeserializeToBusinessModel(rawMsg); // 把Redis原始消息转成业务模型
        observer.OnNext(msg);
    });
    // 注册取消回调:所有消费者都断开时自动退订Redis
    cancelToken.Register(() => redisSub.Unsubscribe());
    return Task.CompletedTask;
})
// 核心:用Publish+RefCount做自动引用计数管理
// 第一个WebSocket接入时自动建立Redis订阅,最后一个WebSocket断开时自动释放Redis连接
.Publish()
.RefCount();

// 每个WebSocket连接建立时,执行这段处理逻辑
async Task ProcessNewWebSocketConnection(WebSocket ws)
{
    // 加发送锁,避免WebSocket多线程并发发送触发报错
    var sendLock = new SemaphoreSlim(1, 1);
    // 订阅全局消息流
    var msgSubscription = redisMessageStream
        .SelectMany(async msg =>
        {
            try
            {
                await sendLock.WaitAsync();
                if (ws.State == WebSocketState.Open)
                {
                    var sendBytes = Encoding.UTF8.GetBytes(JsonSerializer.Serialize(msg));
                    await ws.SendAsync(sendBytes, WebSocketMessageType.Text, true, CancellationToken.None);
                }
            }
            catch (Exception ex)
            {
                // 只记录当前连接的错误,不要把异常抛到外层序列
                Log.Warning(ex, "当前WebSocket推送消息失败,准备断开连接");
            }
            finally
            {
                sendLock.Release();
            }
            return Unit.Default;
        })
        .Subscribe();

    try
    {
        // 维持WebSocket连接,直到客户端主动关闭或者连接异常
        var receiveBuffer = new byte[1024];
        while (ws.State == WebSocketState.Open)
        {
            var receiveResult = await ws.ReceiveAsync(receiveBuffer, CancellationToken.None);
            if (receiveResult.MessageType == WebSocketMessageType.Close)
            {
                await ws.CloseAsync(WebSocketCloseStatus.NormalClosure, "连接正常关闭", CancellationToken.None);
                break;
            }
        }
    }
    finally
    {
        // 连接断开时立刻释放订阅,从广播列表里移除当前消费者,避免内存泄漏
        msgSubscription.Dispose();
        sendLock.Dispose();
    }
}

关键注意点

  • 别用手动调用Connect()的写法:你没法精准控制连接时机,也没法在所有消费者都断开时自动释放上游资源,RefCount()就是专门为这种动态多消费者场景设计的,自动管理引用计数,不需要手动调Connect。
  • 单个消费者的异常一定要在订阅内部处理掉:Rx默认规则是只要一个订阅者抛出未处理异常,整个序列会直接终止,所有订阅者都会被断流,这个坑在多消费者场景下影响极大。
  • WebSocket发送必须加锁:绝大多数WebSocket实现不支持并发Send操作,不加锁会随机出现发送失败、连接意外断开的问题。
  • 如果你需要给新连接的WebSocket重放最近几条消息(比如刚连上先推最新的缓存消息),把.Publish().RefCount()换成.Replay(需要缓存的消息条数).RefCount()即可,会自动缓存指定数量的消息推给新订阅者。

内容的提问来源于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:33:28