FTX WebSocket订单簿校验失败后如何安全重连并清理内存?
问题:FTX订单簿WebSocket连接校验失败后的自动重连实现
我需要通过WebSocket连接FTX交易所,获取BTC/USD订单簿的实时数据流。流程是先获取订单簿快照,之后用WebSocket返回的更新数据重建本地订单簿。每次更新后要通过crc32校验和验证同步状态,如果校验不匹配,就需要重置连接——先取消订阅频道再重新订阅,同时清理全局对象asks、bids、checksum的内存,然后立即重连。
我考虑过在代码末尾加如下死循环来实现重连,但这种方案无法正常停止程序,只能通过关闭终端强制退出,不太理想:
while True: ws.run_forever()
我的现有代码如下:
import websocket,json import zlib from decimal import Decimal import binascii from itertools import chain, zip_longest from typing import Iterable, Sequence asks = {} bids = {} checksum = {'checksum':0} def format_e(dec): return ('{:.' + str(len(dec.as_tuple().digits) - 1) + 'e}').format(dec) def check_sum( asks: Iterable[Sequence[float]], bids: Iterable[Sequence[float]] ) -> int: asks=[[level[0],level[1]]for level in asks.items()] bids=[[level[0],level[1]]for level in bids.items()] order_book_hash_iterator = zip_longest(bids, asks, fillvalue=tuple()) check_string = ":".join( ( str(token) for ask_level, bid_level in order_book_hash_iterator for token in chain(ask_level, bid_level) ) ) return binascii.crc32(check_string.encode("ascii")) def on_open(ws): print('Opened connection') asks.clear() bids.clear() subscribe_message = {'op': 'subscribe', 'channel': 'orderbook','market':'BTC/USD'} ws.send(json.dumps(subscribe_message)) def on_message(ws,message): js=json.loads(message) if js['type'] == 'partial': print('Get Snapshot') for level in js['data']['asks']: asks[level[0]]=level[1] for level in js['data']['bids']: bids[level[0]]=level[1] checksum['checksum']=js['data']['checksum'] if js['type'] == 'update': for level in js['data']['asks']: if level[1]==0: asks.pop(level[0]) else: asks[level[0]]=level[1] for level in js['data']['bids']: if level[1]==0: bids.pop(level[0]) else: bids[level[0]]=level[1] if check_sum(asks,bids) != js['data']['checksum']: print('Error') ws.close() socket = "wss://ftx.com/ws/" ws = websocket.WebSocketApp(socket,on_open=on_open) ws.on_message = lambda ws,msg: on_message(ws,msg) ws.run_forever()
优化方案:优雅重连+可控制停止
我们可以通过添加停止标志和利用on_close回调来实现优雅的重连逻辑,同时支持手动停止程序:
import websocket import json import binascii from itertools import chain, zip_longest from typing import Iterable, Sequence # 全局状态 asks = {} bids = {} checksum = {'checksum': 0} running = True # 控制程序运行的标志 def format_e(dec): return ('{:.' + str(len(dec.as_tuple().digits) - 1) + 'e}').format(dec) def check_sum( asks: Iterable[Sequence[float]], bids: Iterable[Sequence[float]] ) -> int: asks_list = [[level[0], level[1]] for level in asks.items()] bids_list = [[level[0], level[1]] for level in bids.items()] order_book_hash_iterator = zip_longest(bids_list, asks_list, fillvalue=tuple()) check_string = ":".join( ( str(token) for ask_level, bid_level in order_book_hash_iterator for token in chain(ask_level, bid_level) ) ) return binascii.crc32(check_string.encode("ascii")) def on_open(ws): print('已建立连接') # 清理本地数据 asks.clear() bids.clear() checksum['checksum'] = 0 # 订阅订单簿 subscribe_message = {'op': 'subscribe', 'channel': 'orderbook', 'market': 'BTC/USD'} ws.send(json.dumps(subscribe_message)) def on_message(ws, message): global checksum js = json.loads(message) if js['type'] == 'partial': print('获取到订单簿快照') for level in js['data']['asks']: asks[level[0]] = level[1] for level in js['data']['bids']: bids[level[0]] = level[1] checksum['checksum'] = js['data']['checksum'] elif js['type'] == 'update': for level in js['data']['asks']: if level[1] == 0: asks.pop(level[0], None) # 用pop的默认值避免KeyError else: asks[level[0]] = level[1] for level in js['data']['bids']: if level[1] == 0: bids.pop(level[0], None) else: bids[level[0]] = level[1] # 校验逻辑(仅在有校验和时执行) if 'checksum' in js.get('data', {}): local_checksum = check_sum(asks, bids) if local_checksum != js['data']['checksum']: print(f'校验失败:本地校验和{local_checksum},远程校验和{js["data"]["checksum"]}') # 先发送取消订阅消息(可选,但更规范) unsubscribe_msg = {'op': 'unsubscribe', 'channel': 'orderbook', 'market': 'BTC/USD'} try: ws.send(json.dumps(unsubscribe_msg)) except Exception as e: print(f'发送取消订阅消息失败:{e}') # 关闭连接 ws.close() def on_close(ws, close_status_code, close_msg): print(f'连接已关闭,状态码:{close_status_code},消息:{close_msg}') # 如果程序仍在运行,延迟1秒后重连 if running: import time time.sleep(1) create_and_run_ws() def create_and_run_ws(): socket = "wss://ftx.com/ws/" ws = websocket.WebSocketApp( socket, on_open=on_open, on_message=on_message, on_close=on_close ) ws.run_forever() if __name__ == "__main__": try: create_and_run_ws() except KeyboardInterrupt: print('收到停止信号,正在关闭程序') running = False # 可以在这里主动关闭连接 globals().get('ws', None) and ws.close()
关键改进点:
- 可控制停止:添加
running全局标志,按下Ctrl+C触发KeyboardInterrupt时设置running=False,程序会在当前连接关闭后停止重连,优雅退出。 - 优雅重连:在
on_close回调中判断running状态,若为True则延迟1秒后重新创建WebSocket连接并启动,避免频繁重连触发交易所限制。 - 规范的连接重置:校验失败时先尝试发送取消订阅消息,再关闭连接,符合FTX的WebSocket协议规范。
- 鲁棒性提升:使用
pop(level[0], None)避免删除不存在的键时抛出KeyError,添加异常处理避免发送消息失败导致程序崩溃。 - 明确的状态清理:在
on_open中统一清理本地订单簿数据和校验和,确保每次重连后都是干净的初始状态。
内容的提问来源于stack exchange,提问作者apt45
相关产品推荐
相关产品推荐

