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

ASP.NET中基于gRPC流式实现跨方法多客户端通知的问题及优化

用gRPC实现类似SignalR的跨服务消息推送(服务器流式)

问题背景

需要实现类似SignalR中从Hub外部发送消息的功能:让多个客户端通过GUID订阅指定通知器,再从其他方法向订阅该ID的客户端发送消息。要求仅使用gRPC(-web),采用服务器流式方案。

最初实现了LiveService和LiveNotifier类,新消息通知功能正常,但EndLive方法在多客户端场景下无法生效(单客户端时可正常关闭连接)。后续对实现进行了优化修改,现咨询是否有更好的方案或改进建议。

最初实现代码

LiveService

public class LiveService(ILiveNotifier liveNotifier) : LiveServiceBase
{
  private readonly ILiveNotifier _liveNotifier = liveNotifier;

  public override async Task StreamLiveUpdates(
    StreamLiveUpdatesRequest request,
    IServerStreamWriter<string> responseStream,
    ServerCallContext context
  )
  {
    Guid id = GuidMapper.ToGuid(request.Id);
    var cancellationToken = context.CancellationToken;

    while (!cancellationToken.IsCancellationRequested)
    {
      var item = await _liveNotifier.WaitForNewItem(id, cancellationToken);
      await responseStream.WriteAsync(item);
    }
  }
}

新消息通知调用代码

_liveNotifier.NotifyNewItemAvailable(someId, "message to send")

LiveNotifier

public class LiveNotifier : ILiveNotifier
{
  public ConcurrentDictionary<(CancellationToken, Guid), TaskCompletionSource<string>> Waiters { get; set; } = [];

  public void NotifyNewItemAvailable(Guid statusId, string value)
  {
    ConcurrentDictionary<(CancellationToken, Guid), TaskCompletionSource<string>> old = [];
    foreach (var waiter in Waiters)
      if (waiter.Key.Item2 == statusId)
        waiter.Value.TrySetResult(value);
      else
        old.TryAdd(waiter.Key, waiter.Value);

    Waiters = old;
  }

  public void EndLive(Guid statusId)
  {
    ConcurrentDictionary<(CancellationToken, Guid), TaskCompletionSource<string>> old = [];
    foreach (var waiter in Waiters)
      if (waiter.Key.Item2 == statusId) {
        waiter.Value.TrySetCanceled();
      }
      else
        old.TryAdd(waiter.Key, waiter.Value);

    Waiters = old;
  }

  public Task<string> WaitForNewItem(Guid statusId, CancellationToken token)
  {
    TaskCompletionSource<string> newTask = new(TaskCreationOptions.RunContinuationsAsynchronously);
    Waiters.AddOrUpdate((token, statusId), newTask, (_, _) => newTask);

    token.Register(() =>
    {
      newTask.TrySetCanceled(token);
      Waiters.TryRemove((token, statusId), out newTask!);
    });

    return newTask.Task;
  }
}

修改后的实现代码

LiveNotifier

public class LiveNotifier : ILiveNotifier
{
  private readonly ConcurrentDictionary<Guid, ImmutableList<(IServerStreamWriter<string>, TaskCompletionSource)>> _subscriptions = new();

  public Task Subscribe(Guid liveStreamId, IServerStreamWriter<string> responseStream, CancellationToken token)
  {
    TaskCompletionSource tcs = new(TaskCreationOptions.RunContinuationsAsynchronously);
    var newVal = (responseStream, tcs);

    _subscriptions.AddOrUpdate(liveStreamId, [newVal], (key, oldValue) => oldValue.Add(newVal));

    token.Register(
      () =>
      {
        tcs.TrySetCanceled();
        _subscriptions.AddOrUpdate(liveStreamId, [], (id, oldValues) => oldValues.Remove(newVal));
      },
      false
    );
    return tcs.Task;
  }

  public void Broadcast(Guid liveStatusId, string message)
  {
    _subscriptions.TryGetValue(
      liveStatusId,
      out ImmutableList<(IServerStreamWriter<string>, TaskCompletionSource)>? values
    );
    if (values != null)
    {
      foreach (var req in values)
        req.Item1.WriteAsync(message);
    }
  }

  public void EndLive(Guid liveStatusId)
  {
    _subscriptions.TryGetValue(
      liveStatusId,
      out ImmutableList<(IServerStreamWriter<string>, TaskCompletionSource)>? values
    );
    if (values != null)
    {
      foreach (var req in values)
        req.Item2.TrySetResult();
      _subscriptions.TryRemove(liveStatusId, out _);
    }
  }
}

LiveService代码片段

Guid id = GuidMapper.ToGuid(request.Id);
var cancellationToken = context.CancellationToken;

await _liveNotifier.Subscribe(id, responseStream, cancellationToken);

改进建议

1. 处理异步消息发送的等待与异常

修改后的Broadcast方法直接调用WriteAsync但未等待,会导致未处理的异步任务,可能引发异常或消息顺序混乱。建议改为异步实现并捕获发送失败的情况,自动清理无效订阅:

public async Task Broadcast(Guid liveStatusId, string message)
{
    if (_subscriptions.TryGetValue(liveStatusId, out var values))
    {
        foreach (var req in values)
        {
            try
            {
                await req.Item1.WriteAsync(message);
            }
            catch (Exception ex)
            {
                // 移除发送失败的无效订阅
                _subscriptions.AddOrUpdate(liveStatusId, [], (id, oldValues) => oldValues.Remove(req));
            }
        }
    }
}

2. 优化订阅清理的并发安全性

在Subscribe的取消回调中,使用AddOrUpdate移除订阅时,可能存在并发冲突。建议确保ImmutableList的原子更新,同时避免空列表残留:

token.Register(() =>
{
    tcs.TrySetCanceled();
    _subscriptions.AddOrUpdate(liveStreamId, ImmutableList<(IServerStreamWriter<string>, TaskCompletionSource)>.Empty, 
        (id, oldValues) => 
        {
            var updated = oldValues.Remove(newVal);
            return updated.IsEmpty ? ImmutableList<(IServerStreamWriter<string>, TaskCompletionSource)>.Empty : updated;
        });
}, false);

3. 明确EndLive的语义与客户端通知

当前EndLive仅完成Task让连接关闭,建议先给客户端发送明确的结束标记,再终止连接,让客户端区分主动结束与异常断开:

public async Task EndLive(Guid liveStatusId)
{
    if (_subscriptions.TryRemove(liveStatusId, out var values))
    {
        foreach (var req in values)
        {
            try
            {
                await req.Item1.WriteAsync("STREAM_TERMINATED"); // 自定义结束标记
            }
            catch { }
            req.Item2.TrySetResult();
        }
    }
}

4. 权衡数据结构的性能与安全性

ImmutableList每次更新都会创建新实例,高并发场景下可能有性能开销。若订阅/取消频率极高,可改用ConcurrentDictionary<Guid, ConcurrentBag<(IServerStreamWriter<string>, TaskCompletionSource)>>,并定期清理空Bag;若频率较低,ImmutableList的线程安全特性更省心。

5. 添加关键节点日志监控

在订阅、消息发送、流结束等环节添加日志,方便排查问题:

  • 订阅时记录GUID与客户端标识(从ServerCallContext获取)
  • 消息发送失败时记录异常信息
  • EndLive时记录GUID与订阅数量

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 08:49:57