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

Reactive.NET处理Binance复合WebSocket流拆分输出指定IObservable类型问询

Binance复合WebSocket流拆分实现需求

现有Binance复合WebSocket流地址:wss://stream.binance.com:9443/stream?streams=bnbusdt@ticker/dogeusdt@depth5,需要实现以下两个可观测序列属性:

public IObservable<WebSocketPriceTicker24Hr> Tickers => ...;
public IObservable<WebSocketDepth> Depth => ...;

流日志示例

Connection opened
Message: {"stream":"dogeusdt@depth5","data":{"lastUpdateId":3740272226,"bids":[["0.20140000","21189.00000000"],["0.20130000","275878.00000000"],["0.20120000","290900.00000000"],["0.20110000","313592.00000000"],["0.20100000","367368.00000000"]],"asks":[["0.20150000","109090.00000000"],["0.20160000","404515.00000000"],["0.20170000","649409.00000000"],["0.20180000","360650.00000000"],["0.20190000","185381.00000000"]]}}
Message: {"stream":"bnbusdt@ticker","data":{"e":"24hrTicker","E":1638097890123,"s":"BNBUSDT","p":"-2.50000000","P":"-0.416","w":"598.07225116","x":"601.40000000","c":"599.00000000","Q":"0.45200000","b":"599.00000000","B":"122.06600000","a":"599.10000000","A":"0.54000000","o":"601.50000000","h":"621.30000000","l":"572.40000000","v":"1286613.77200000","q":"769487994.99120000","O":1638011490067,"C":1638097890067,"F":471394573,"L":472263211,"n":868639}}

问题说明

流返回的消息需要先反序列化为WebSocketResponse<T>类型,再按消息类型拆分,最终输出对应类型的可观测序列。现有已完成代码如下:

public IObservable<string> Messages => Observable
    .FromEventPattern<MessageReceivedEventArgs>(h => _webSocket.MessageReceived += h,
        h => _webSocket.MessageReceived -= h)
    .Select(e => e.EventArgs.Message);

// 实体模型定义
public class WebSocketResponse<T>
{
    public string? Stream { get; set; }
    public T? Data { get; set; }
}

public class WebSocketPriceTicker24Hr
{
    [JsonPropertyName("e")] public string? EventType { get; set; }
    [JsonPropertyName("E")] public long EventTime { get; set; }
    [JsonPropertyName("s")] public string? Symbol { get; set; }
    [JsonPropertyName("p")] public decimal PriceChange { get; set; }
    [JsonPropertyName("P")] public decimal PriceChangePercent { get; set; }
    [JsonPropertyName("w")] public decimal WeightedAveragePrice { get; set; }
    [JsonPropertyName("x")] public decimal PreviousClosePrice { get; set; }
    [JsonPropertyName("c")] public decimal LastPrice { get; set; }
    [JsonPropertyName("Q")] public decimal LastQuantity { get; set; }
    [JsonPropertyName("b")] public decimal BestBidPrice { get; set; }
    [JsonPropertyName("B")] public decimal BestBidQuantity { get; set; }
    [JsonPropertyName("a")] public decimal BestAskPrice { get; set; }
    [JsonPropertyName("A")] public decimal BestAskQuantity { get; set; }
    [JsonPropertyName("o")] public decimal OpenPrice { get; set; }
    [JsonPropertyName("h")] public decimal HighPrice { get; set; }
    [JsonPropertyName("l")] public decimal LowPrice { get; set; }
    [JsonPropertyName("v")] public decimal TotalTradedBaseVolume { get; set; }
    [JsonPropertyName("q")] public decimal TotalTradedQuoteVolume { get; set; }
    [JsonPropertyName("O")] public long OpenTime { get; set; }
    [JsonPropertyName("C")] public long CloseTime { get; set; }
    [JsonPropertyName("F")] public long FirstTradeId { get; set; }
    [JsonPropertyName("L")] public long LastTradeId { get; set; }
    [JsonPropertyName("n")] public long Count { get; set; }
}

实现步骤

1. 修正WebSocketDepth实体映射

现有WebSocketDepth类的字段映射和Binance深度快照接口返回字段不匹配,调整为:

public class WebSocketDepth
{
    [JsonPropertyName("lastUpdateId")] 
    public long LastUpdateId { get; set; }

    [JsonPropertyName("bids")]
    public IEnumerable<IEnumerable<string>> Bids { get; set; } = Array.Empty<IEnumerable<string>>();

    [JsonPropertyName("asks")]
    public IEnumerable<IEnumerable<string>> Asks { get; set; } = Array.Empty<IEnumerable<string>>();
}

2. 实现两个可观测序列属性

基于已有的Messages流直接拆分,不需要在全局订阅里处理反序列化逻辑,符合Rx的流拆分设计:

// 类字段:复用Json序列化配置,避免重复开销
private readonly JsonSerializerOptions _jsonSerializerOptions = new JsonSerializerOptions
{
    PropertyNameCaseInsensitive = true,
    NumberHandling = JsonNumberHandling.AllowReadingFromString
};

public IObservable<WebSocketPriceTicker24Hr> Tickers => Messages
    .Select(msg =>
    {
        try
        {
            return JsonSerializer.Deserialize<WebSocketResponse<WebSocketPriceTicker24Hr>>(msg, _jsonSerializerOptions);
        }
        catch
        {
            return null;
        }
    })
    .Where(res => res != null && res.Stream?.EndsWith("@ticker") == true && res.Data != null)
    .Select(res => res.Data!);

public IObservable<WebSocketDepth> Depth => Messages
    .Select(msg =>
    {
        try
        {
            return JsonSerializer.Deserialize<WebSocketResponse<WebSocketDepth>>(msg, _jsonSerializerOptions);
        }
        catch
        {
            return null;
        }
    })
    .Where(res => res != null && res.Stream?.Contains("@depth") == true && res.Data != null)
    .Select(res => res.Data!);

3. 订阅示例

直接订阅两个拆分后的序列即可:

// 用CompositeDisposable管理多个订阅
_eventSubscription = new CompositeDisposable(
    // 订阅24小时行情流
    Tickers.Subscribe(ticker =>
    {
        Console.WriteLine($"[{ticker.Symbol}] 最新价:{ticker.LastPrice}, 涨跌幅:{ticker.PriceChangePercent}%");
    }),
    // 订阅深度流
    Depth.Subscribe(depth =>
    {
        var bestBid = depth.Bids.First();
        var bestAsk = depth.Asks.First();
        Console.WriteLine($"深度更新 | 买一价:{bestBid.First()} 数量:{bestBid.Last()} | 卖一价:{bestAsk.First()} 数量:{bestAsk.Last()}");
    })
);

原来的全局Messages订阅里的反序列化逻辑可以直接删除。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 20:45:10