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

如何让子任务触发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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 17:45:32