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

