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

Binance API深度订单簿Websocket数据处理延迟与ID异常问题

开发Binance API Websocket连接时的订单簿同步问题

我在开发Binance API的Websocket连接时,遇到了处理延迟过高与数据溢出问题。当前采用双进程架构:WebsocketProcessor负责建立Websocket连接并推送数据至数据库,DataProcessor负责从数据库拉取数据并处理。核心问题是无法满足本地订单簿更新的官方规则,尤其是规则5(首个处理事件需满足U ≤ lastUpdateId+1 且 u ≥ lastUpdateId+1)始终不成立,同时UpdateID链出现断层,进入orderBookUpdate循环时会触发约10条更新缺失警告。

相关日志信息

INFO:root:Websocket connection established
INFO:root:Initiated trades: {'_id': 1, 'trades': []}
INFO:root:Deleted 6 updates from depth bufferINFO:root:Order book initialization updateId not stable
INFO:root:Initiated orderbook: ...
INFO:root:Missing depth update id. Current id: 47878509997, next id: 47878507790
INFO:root:Missing depth update id. Current id: 47878507790, next id: 47878507809
INFO:root:Missing depth update id. Current id: 47878507809, next id: 47878507822
INFO:root:Missing depth update id. Current id: 47878507822, next id: 47878507832
INFO:root:Missing depth update id. Current id: 47878507832, next id: 47878507839
INFO:root:Missing depth update id. Current id: 47878507839, next id: 47878507847

DataProcessor中订单簿初始化与更新的代码

def initiateOrderBook(self):
    # 检查数据库中是否已有更新数据
    while True:
        # 清空旧更新并取出当前数据
        updates = self.collection.find_one_and_update(
            filter={"_id": 4},
            update={
                "$set": {"depthUpdates": []}  # 将depthUpdates重置为空列表
            },
            projection={"depthUpdates": 1},  # 仅返回取出的更新数据
            return_document=pymongo.ReturnDocument.BEFORE
        )

        if updates is None:
            logging.info('数据库中暂无更新数据,重试中...')
            time.sleep(2)
        else:
            break

    # 按updateId排序更新数据
    updates = SortedDict(list((x['u'], x) for x in updates['depthUpdates']))

    # 获取订单簿快照
    orderBook = self.getOrderbook()

    # 过滤掉快照前的旧更新
    c = 0
    for last_update_id, data in updates.items():
        if last_update_id <= orderBook['lastUpdateId']:
            updates.pop(last_update_id)
            c += 1
    logging.info(f'从深度缓冲区中删除了{c}条旧更新')

    # 验证订单簿与更新数据的连续性
    first_item = updates.peekitem()[1]
    if first_item['U'] <= orderBook['lastUpdateId'] + 1 <= first_item['u']:
        logging.info("订单簿已准备好接收更新")
    else:
        logging.info("订单簿初始化阶段UpdateID不连续")

    # 将剩余更新数据推回数据库
    filter_criteria = {"_id": 4}
    update = {
        "$push": {"depthUpdates": {"$each": list(updates.values())}}
    }

    self.collection.update_one(filter_criteria, update)

    # 将订单簿快照存入数据库
    document = self.collection.find_one_and_update(
        filter={"_id": 2},
        update={
            "$set": {"orderBook": orderBook}
        },
        upsert=True,
        return_document=pymongo.ReturnDocument.AFTER
    )
    logging.info(f'已初始化订单簿: {document}')

def orderBookUpdate(self, updatingSpeed):
    while True:
        # 从数据库获取当前订单簿
        document = self.collection.find_one(filter={"_id": 2})

        # 取出更新数据并清空缓冲区
        updates = self.collection.find_one_and_update(
            filter={"_id": 4},
            update={
                "$set": {"depthUpdates": []}  # 将depthUpdates重置为空列表
            },
            projection={"depthUpdates": 1},  # 仅返回取出的更新数据
            return_document=pymongo.ReturnDocument.BEFORE
        )

        # 处理更新数据
        orderBook = document['orderBook']
        updates = SortedDict(list((x['u'], x) for x in updates['depthUpdates']))  # 按updateId排序

        # 订单簿更新逻辑
        for entry in updates.values():
            # 检查UpdateID链的完整性
            if orderBook['lastUpdateId'] + 1 != entry['u']:
                logging.info(f'缺失深度更新ID. 当前ID: {orderBook["lastUpdateId"]}, 下一个ID: {entry["u"]}')

            # 更新买盘和卖盘数据
            for side in ['b', 'a']:
                for price_level in entry[side]:
                    quantity = float(price_level[1])
                    if quantity == 0:
                        orderBook[side].pop(price_level[0], None)
                    else:
                        orderBook[side][price_level[0]] = quantity

            # 更新当前订单簿的lastUpdateId
            orderBook['lastUpdateId'] = entry['u']

        # 限制订单簿深度为最近的1000条数据(买盘500条,卖盘500条)
        orderBook['b'] = SortedDict(list(orderBook['b'].items())[-500:])
        orderBook['a'] = SortedDict(list(orderBook['a'].items())[:500])

        # 将更新后的订单簿推回数据库
        self.collection.update_one(filter={"_id": 2}, update={"$set": {"orderBook": orderBook}})

        # 控制更新频率
        time.sleep(updatingSpeed)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 15:16:01