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

如何在Observable内创建新Observable以单独获取OrderBook增量更新数据

解决方案

1. 优化订阅复用(避免重复WebSocket连接)

首先添加通道缓存,确保同一个Instrument只订阅一次WebSocket通道,减少资源消耗:

private final Map<Instrument, Observable<JsonNode>> channelCache = new ConcurrentHashMap<>();
private final Map<Instrument, OrderBook> orderBookMap = new HashMap<>();

private Observable<JsonNode> getChannel(Instrument instrument) {
    return channelCache.computeIfAbsent(instrument, service::subscribeChannel);
}

2. 保留原getOrderBook方法(返回完整OrderBook)

修改原方法使用共享通道,同时完善空值场景处理:

public Observable<OrderBook> getOrderBook(Instrument instrument) {
    return getChannel(instrument).flatMap(jsonNode -> {
        String action = jsonNode.get("action").asText().toLowerCase();
        if ("snapshot".equals(action)) {
            // 解析完整快照并缓存
            OrderBook orderBook = mapper.treeToValue(
                    jsonNode.get("data"),
                    OrderBook.class
            );
            orderBookMap.put(instrument, orderBook);
            return Observable.just(orderBook);
        } else {
            OrderBook orderBook = orderBookMap.get(instrument);
            if (orderBook == null) {
                // 未收到快照时忽略增量更新,可根据业务需求调整逻辑
                return Observable.empty();
            }
            List<PublicOrder> incrementalUpdateData = mapper.treeToValue(
                    jsonNode.get("data").get(0).get("asks"),
                    mapper.getTypeFactory().constructCollectionType(List.class, PublicOrder.class)
            );
            orderBook.update(incrementalUpdateData);
            return Observable.just(orderBook);
        }
    });
}

3. 新增getOrderBookUpdate方法(仅返回增量数据)

过滤快照消息,只解析并发射增量更新内容:

public Observable<List<PublicOrder>> getOrderBookUpdate(Instrument instrument) {
    return getChannel(instrument)
            // 只处理增量更新类型的消息
            .filter(jsonNode -> !"snapshot".equalsIgnoreCase(jsonNode.get("action").asText()))
            // 解析出增量更新数据
            .map(jsonNode -> mapper.treeToValue(
                    jsonNode.get("data").get(0).get("asks"),
                    mapper.getTypeFactory().constructCollectionType(List.class, PublicOrder.class)
            ))
            // 确保已初始化订单簿后再发射增量,避免无效数据
            .filter(update -> orderBookMap.containsKey(instrument));
}

使用方式

  • 订阅完整OrderBook:
getOrderBook(someInstrument).subscribe(orderBook -> {
    // 处理完整订单簿数据
});
  • 订阅增量更新数据:
getOrderBookUpdate(someInstrument).subscribe(incrementalData -> {
    // 仅处理增量更新的订单数据
});

内容的提问来源于stack exchange,提问作者Илья Смирнов

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 02:15:34