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

如何改造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("客户端正常断开");
    }
}

核心改造点说明

  1. 多客户端订阅支持

    • 使用ConcurrentBag<ChannelWriter<Notification>>维护所有订阅客户端的写入器,每个客户端订阅时创建独立的Channel,保证消息能推送给所有在线客户端。
  2. 客户端断开检测与日志

    • 利用CallContext的CancellationToken:客户端断开时该令牌会触发取消,通过Register回调完成当前客户端的Channel.Writer,并从订阅集合中移除,同时记录断开日志。
    • 写入通知时的兜底处理:如果TryWrite失败或抛出异常,说明客户端已异常断开,同样清理订阅者并记录日志,避免无效订阅者占用资源。
  3. 线程安全与异常隔离

    • 用ConcurrentBag保证多线程下集合操作的安全性;每个订阅者的写入操作单独捕获异常,防止单个客户端的异常导致整个通知流程中断。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 03:55:12