如何让子任务触发WebSocket的SendAsync方法高效推送实时数据?
可行技术方案
针对你的需求,推荐以下几种无需轮询、基于事件/消息触发的WebSocket推送方案,均能实现子任务更新时直接触发推送:
1. 进程内发布-订阅(Pub/Sub)事件总线
这是最通用、可扩展的方案,适合后续继续新增数据类型的场景:
- 定义统一的
DataUpdateEvent模型,包含数据类型标识、数据内容等核心字段。 - 实现轻量级内存事件总线,提供订阅、发布接口,无需依赖第三方MQ(进程内通信足够高效)。
- 各子任务作为发布者,在数据更新时发布对应事件;WebSocket连接管理模块作为订阅者,监听所有需要推送的事件类型,收到事件后直接调用
SendAsync推送。
代码示例(C#):
// 统一事件模型 public class DataUpdateEvent { public string DataType { get; set; } // 如"StockPrice"/"Position"/"Order" public object Data { get; set; } } // 轻量级内存事件总线 public class InMemoryEventBus { private readonly Dictionary<string, List<Action<DataUpdateEvent>>> _subscribers = new(); public void Subscribe(string dataType, Action<DataUpdateEvent> handler) { if (!_subscribers.ContainsKey(dataType)) _subscribers[dataType] = new List<Action<DataUpdateEvent>>(); _subscribers[dataType].Add(handler); } public void Publish(DataUpdateEvent @event) { if (_subscribers.TryGetValue(@event.DataType, out var handlers)) { foreach (var handler in handlers) handler.Invoke(@event); } } } // WebSocket管理类订阅事件并推送 public class WebSocketHandler { private readonly InMemoryEventBus _eventBus; private readonly WebSocket _webSocket; public WebSocketHandler(InMemoryEventBus eventBus, WebSocket webSocket) { _eventBus = eventBus; _webSocket = webSocket; // 订阅所有需要推送的数据类型 _eventBus.Subscribe("StockPrice", PushUpdate); _eventBus.Subscribe("Position", PushUpdate); _eventBus.Subscribe("Order", PushUpdate); } private async void PushUpdate(DataUpdateEvent @event) { if (_webSocket.State != WebSocketState.Open) return; var json = JsonSerializer.Serialize(@event); var buffer = Encoding.UTF8.GetBytes(json); await _webSocket.SendAsync( new ArraySegment<byte>(buffer), WebSocketMessageType.Text, endOfMessage: true, CancellationToken.None ); } } // 子任务发布事件(以股价更新为例) public class StockPriceMonitor { private readonly InMemoryEventBus _eventBus; public StockPriceMonitor(InMemoryEventBus eventBus) { _eventBus = eventBus; } public void OnPriceChanged(StockPrice newPrice) { _eventBus.Publish(new DataUpdateEvent { DataType = "StockPrice", Data = newPrice }); } }
2. 直接注入推送委托(轻量场景首选)
如果系统复杂度较低,不需要通用事件总线,可以直接给子任务注入推送委托,省去事件总线的中间层:
- 定义推送委托类型,封装WebSocket的
SendAsync逻辑。 - 子任务初始化时传入该委托,数据更新时直接调用委托触发推送。
代码示例:
// 定义推送委托 public delegate Task PushDataDelegate(string dataType, object data); // 仓位更新子任务 public class PositionManager { private readonly PushDataDelegate _pushDelegate; public PositionManager(PushDataDelegate pushDelegate) { _pushDelegate = pushDelegate; } public async void OnPositionUpdated(Position updatedPosition) { await _pushDelegate("Position", updatedPosition); } } // WebSocket端初始化委托并注入子任务 var webSocket = /* 获取已建立的WebSocket连接 */; var pushDelegate = new PushDataDelegate(async (dataType, data) => { if (webSocket.State != WebSocketState.Open) return; var payload = JsonSerializer.Serialize(new { DataType = dataType, Data = data }); var buffer = Encoding.UTF8.GetBytes(payload); await webSocket.SendAsync( new ArraySegment<byte>(buffer), WebSocketMessageType.Text, endOfMessage: true, CancellationToken.None ); }); var positionManager = new PositionManager(pushDelegate);
3. Channel生产者-消费者模式(高并发场景)
如果子任务更新频率极高,推荐使用.NET的Channel(或其他语言的类似组件)实现生产者-消费者模式,避免并发推送冲突:
- 创建有界
Channel,控制队列容量防止内存溢出。 - 子任务作为生产者,将更新数据写入Channel;WebSocket作为消费者,持续从Channel读取数据并推送。
代码示例:
// 创建有界Channel(最多缓存100条消息,超出时丢弃最旧的) var updateChannel = Channel.CreateBounded<DataUpdateEvent>(new BoundedChannelOptions(100) { FullMode = BoundedChannelFullMode.DropOldest }); // 订单更新子任务(生产者) public class OrderProcessor { private readonly ChannelWriter<DataUpdateEvent> _channelWriter; public OrderProcessor(ChannelWriter<DataUpdateEvent> channelWriter) { _channelWriter = channelWriter; } public async void OnOrderUpdated(Order updatedOrder) { await _channelWriter.WriteAsync(new DataUpdateEvent { DataType = "Order", Data = updatedOrder }); } } // WebSocket消费者任务 async Task StartWebSocketConsumer(WebSocket webSocket, ChannelReader<DataUpdateEvent> channelReader) { while (await channelReader.WaitToReadAsync()) { if (channelReader.TryRead(out var updateEvent) && webSocket.State == WebSocketState.Open) { var json = JsonSerializer.Serialize(updateEvent); var buffer = Encoding.UTF8.GetBytes(json); await webSocket.SendAsync( new ArraySegment<byte>(buffer), WebSocketMessageType.Text, endOfMessage: true, CancellationToken.None ); } } } // 启动消费者任务 _ = StartWebSocketConsumer(webSocket, updateChannel.Reader);
关键注意事项
- 连接状态检查:推送前务必检查
WebSocket.State,避免在连接已关闭时调用SendAsync抛出异常。 - 线程安全:如果多个子任务同时推送,可通过锁或Channel的线程安全特性避免并发冲突。
- 序列化统一:统一使用JSON或其他格式序列化数据,并携带数据类型标识,方便客户端解析。
内容的提问来源于stack exchange,提问作者Maya
相关产品推荐
相关产品推荐

