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
相关产品推荐
相关产品推荐

