如何使用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
相关产品推荐
相关产品推荐

