如何改造protobuf-net.grpc通知系统以支持多客户端订阅及断开日志
改造protobuf-net.grpc服务端流实现多客户端订阅与断开日志
问题根源
原代码使用单个Channel实现生产者/消费者模式,同一时刻仅能有一个客户端消费消息,无法满足多客户端订阅需求;同时未监听客户端断开的取消信号,无法记录断开事件。
改造后完整代码
private readonly INotificationService _notificationService; // 线程安全集合,维护所有订阅客户端的写入器 private readonly ConcurrentBag<ChannelWriter<Notification>> _subscribers = new ConcurrentBag<ChannelWriter<Notification>>(); public ClientNotificationService(INotificationService notificationService) { _notificationService = notificationService; _notificationService.OnNotification += OnNotification; } private async void OnNotification(object sender, Notification notification) { // 遍历所有订阅者发送通知,单独捕获异常避免单个客户端问题影响全局 var tasks = _subscribers.Select(async writer => { try { if (!writer.TryWrite(notification)) { // 写入失败说明客户端已断开,移除订阅者并记录 _subscribers.TryTake(out writer); LogClientDisconnect(); } } catch (Exception ex) { // 处理写入异常,清理订阅者并记录日志 _subscribers.TryTake(out writer); LogClientDisconnect(ex); } }); await Task.WhenAll(tasks); } public async IAsyncEnumerable<Notification> SubscribeAsync([EnumeratorCancellation] CallContext context = default) { // 为每个订阅客户端创建独立的Channel var channel = Channel.CreateUnbounded<Notification>(); _subscribers.Add(channel.Writer); // 监听客户端断开的取消令牌,触发时清理资源 context.CancellationToken.Register(() => { channel.Writer.Complete(); if (_subscribers.TryTake(out var writer) && writer == channel.Writer) { LogClientDisconnect(); } }); // 返回当前客户端的消息流 await foreach (var notification in channel.Reader.ReadAllAsync(context.CancellationToken)) { yield return notification; } } // 自定义日志方法,根据实际业务需求替换为日志框架调用 private void LogClientDisconnect(Exception ex = null) { if (ex != null) { Console.WriteLine($"客户端异常断开:{ex.Message}"); } else { Console.WriteLine("客户端正常断开"); } }
核心改造点说明
多客户端订阅支持
- 使用
ConcurrentBag<ChannelWriter<Notification>>维护所有订阅客户端的写入器,每个客户端订阅时创建独立的Channel,保证消息能推送给所有在线客户端。
- 使用
客户端断开检测与日志
- 利用
CallContext的CancellationToken:客户端断开时该令牌会触发取消,通过Register回调完成当前客户端的Channel.Writer,并从订阅集合中移除,同时记录断开日志。 - 写入通知时的兜底处理:如果
TryWrite失败或抛出异常,说明客户端已异常断开,同样清理订阅者并记录日志,避免无效订阅者占用资源。
- 利用
线程安全与异常隔离
- 用
ConcurrentBag保证多线程下集合操作的安全性;每个订阅者的写入操作单独捕获异常,防止单个客户端的异常导致整个通知流程中断。
- 用
内容的提问来源于stack exchange,提问作者user2727133
相关产品推荐
相关产品推荐

