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

