Binance期货Websocket数据流订单取消异常排查及代码审核请求
问题描述
我正在用Binance期货的UserData Stream实现订单取消逻辑:当止盈(TP)触发时取消对应的止损(SL)订单,止损触发时取消对应的止盈订单。但目前存在部分TP/SL触发后,反向订单未被取消的问题。由于Binance期货没有现货的OCO(一键平仓对冲)功能,我自行编写了代码实现该逻辑,烦请帮忙检查代码是否存在错误:
import datetime import json import threading import requests import websocket import sys import ccxt import random from binance.client import Client Api = 'None' Secret = 'None' BinanceFutures = ccxt.binance({'apiKey': Api,'secret': Secret,'enableRateLimit': True,'options': {'defaultType': 'future','defaultMarket':'futures'}}) BinanceFutures.load_markets() client = Client(Api,Secret) SummarySavingList = [] ProblemsSavingList = [] def ProblemSaver(): ProblemsSavedList = [] while True: for Problem in ProblemsSavingList: if not Problem in ProblemsSavedList: Symbol = Problem.split(',')[0].split(":")[1] ProblemName = Problem.split(',')[1].split(":")[1] Where = Problem.split(',')[2].split(":")[1] restorePoint = sys.stdout sys.stdout = sys.stdout sys.stdout = open(f"C:\\Users\\mujah\\OneDrive\\Desktop\\RobotTrials\\Problems.txt","a") print(f"Message:{ProblemName} from {Where} when running functions on {Symbol}......") sys.stdout.close() sys.stdout = restorePoint ProblemsSavedList.append(Problem) def PositionCloser(Coin): try: PositionInfo = client.futures_position_information(symbol=Coin) if float(PositionInfo[0]['positionAmt']) != 0: SELL = client.futures_create_order(symbol=Coin,side='SELL',type='MARKET',quantity=abs(float(PositionInfo[0]['positionAmt']))) except Exception as e: print(e) def SummarySaver(): SummarySavedList = [] while True: for Summary in SummarySavingList: if not Summary in SummarySavedList: Symbol = Summary.split(',')[0].split(":")[1] Side = Summary.split(',')[1].split(":")[1] Result = Summary.split(',')[2].split(":")[1] SellTiming = Summary.split(',')[3].split(":")[1] restorePoint = sys.stdout sys.stdout = sys.stdout sys.stdout = open(f"C:\\Users\\mujah\\OneDrive\\Desktop\\RobotTrials\\Summary.txt","a") print(f"{Result} Hit When {Side} Trade is Taken On {Symbol},It Was Sold on {SellTiming}...") sys.stdout.close() sys.stdout = restorePoint SummarySavedList.append(Summary) def StartStreaming(): try: BINANCE_FUTURES_END_POINT = "https://fapi.binance.com/fapi/v1/listenKey" def create_futures_listen_key(api_key): response = requests.post(url=BINANCE_FUTURES_END_POINT, headers={'X-MBX-APIKEY': api_key}) return response.json()['listenKey'] FUTURES_STREAM_END_POINT_1 = "wss://fstream.binance.com" listen_key = create_futures_listen_key(Api) futures_connection_url = f"{FUTURES_STREAM_END_POINT_1}/ws/{listen_key}" listen_key = create_futures_listen_key(Api) url = f"{FUTURES_STREAM_END_POINT_1}/ws/{listen_key}" def on_open(ws): print(f"Open: futures order stream connected") def on_message(ws, message): message = json.loads(message) print(f"Message: {message}") if message["e"] == 'ORDER_TRADE_UPDATE': Symbol = message["o"]["s"] Type = message["o"]["o"] Status = message["o"]["x"] Side = message["o"]['S'] if Status == 'FILLED' and Side == 'SELL': if Type == 'TAKE_PROFIT_MARKET': try: Cancel = client.futures_cancel_all_open_orders(symbol=Symbol) except Exception as e: print(e) SellTiming = datetime.datetime.now().strftime("%Y-%m-%d %H-%M-%S") PositionCloser(Coin=Symbol) SummarySavingList.append(f"Symbol:{Symbol},Side:LONG,Result:Takeprofit,SellTiming:{SellTiming}") elif Type == 'STOP_MARKET' and Side == 'SELL': try: Cancel = client.futures_cancel_all_open_orders(symbol=Symbol) except Exception as e: print(e) SellTiming = datetime.datetime.now().strftime("%Y-%m-%d %H-%M-%S") PositionCloser(Coin=Symbol) SummarySavingList.append(f"Symbol:{Symbol},Side:LONG,Result:StopLoss,SellTiming:{SellTiming}") elif Status == 'EXPIRED': if Type == 'TAKE_PROFIT_MARKET' and Side == 'SELL': try: Cancel = client.futures_cancel_all_open_orders(symbol=Symbol) except Exception as e: print(e) SellTiming = datetime.datetime.now().strftime("%Y-%m-%d %H-%M-%S") PositionCloser(Coin=Symbol) SummarySavingList.append(f"Symbol:{Symbol},Side:LONG,Result:TakeProfit,SellTiming:{SellTiming}") elif Type == 'STOP_MARKET' and Side == 'SELL': try: Cancel = client.futures_cancel_all_open_orders(symbol=Symbol) except Exception as e: print(e) SellTiming = datetime.datetime.now().strftime("%Y-%m-%d %H-%M-%S") PositionCloser(Coin=Symbol) SummarySavingList.append(f"Symbol:{Symbol},Side:LONG,Result:StopLoss,SellTiming:{SellTiming}") elif Status == 'PARTIALLY_FILLED': if Type == 'TAKE_PROFIT_MARKET' and Side == 'SELL': try: Cancel = client.futures_cancel_all_open_orders(symbol=Symbol) except Exception as e: print(e) SellTiming = datetime.datetime.now().strftime("%Y-%m-%d %H-%M-%S") PositionCloser(Coin=Symbol) SummarySavingList.append( f"Symbol:{Symbol},Side:LONG,Result:TakeProfit,SellTiming:{SellTiming}") elif Type == 'STOP_MARKET' and Side == 'SELL': try: Cancel = client.futures_cancel_all_open_orders(symbol=Symbol) except Exception as e: print(e) SellTiming = datetime.datetime.now().strftime("%Y-%m-%d %H-%M-%S") PositionCloser(Coin=Symbol) SummarySavingList.append( f"Symbol:{Symbol},Side:LONG,Result:StopLoss,SellTiming:{SellTiming}") def on_error(ws, error): print(f"Error: {error}") def on_close(ws, close_status_code, close_msg): print(f"Close: {close_status_code} {close_msg}") def stream_ticker(): ws = websocket.WebSocketApp(url=url,on_open=on_open,on_message=on_message,on_error=on_error,on_close=on_close) ws.run_forever() stream_ticker() except Exception as e: ProblemsSavingList.append(f"Symbol:None,Problem:{e},Where:Streamer,Number:{random.randint(1,1000)}") pass SummaryThread = threading.Thread(target=SummarySaver) SummaryThread.start() ProblemThread = threading.Thread(target=ProblemSaver) ProblemThread.start() StartStreaming()
代码问题及修复方案
1. 监听Key重复创建+缺失续期逻辑
你连续两次调用create_futures_listen_key(Api),第二次会覆盖第一次生成的Key,导致第一次的连接无效;同时Binance要求listenKey每60分钟续期一次,否则会自动失效,代码中没有处理续期,会导致流连接中断。
修复:
# 在StartStreaming函数中添加续期线程 import time def keep_alive_listen_key(): while True: time.sleep(30 * 60) # 每30分钟续期一次,提前避免过期 try: requests.put(BINANCE_FUTURES_END_POINT, headers={'X-MBX-APIKEY': Api}) print("ListenKey续期成功") except Exception as e: print(f"ListenKey续期失败: {e}") # 修改listenKey创建逻辑,只生成一次 listen_key = create_futures_listen_key(Api) url = f"{FUTURES_STREAM_END_POINT_1}/ws/{listen_key}" # 启动续期线程 keep_alive_thread = threading.Thread(target=keep_alive_listen_key, daemon=True) keep_alive_thread.start()
2. 订单状态判断逻辑嵌套错误
PARTIALLY_FILLED的判断被嵌套在EXPIRED分支内部,永远不会被触发,因为状态不可能同时是EXPIRED和PARTIALLY_FILLED。
修复:
将PARTIALLY_FILLED判断移到外层,与FILLED、EXPIRED同级:
if message["e"] == 'ORDER_TRADE_UPDATE': Symbol = message["o"]["s"] Type = message["o"]["o"] Status = message["o"]["x"] Side = message["o"]['S'] if Status == 'FILLED' and Side == 'SELL': # 原FILLED逻辑不变 elif Status == 'EXPIRED' and Side == 'SELL': # 原EXPIRED逻辑不变 elif Status == 'PARTIALLY_FILLED' and Side == 'SELL': # 原PARTIALLY_FILLED逻辑不变
3. 多线程读写全局列表无锁
SummarySavingList和ProblemsSavingList是全局列表,多线程读写时没有加锁,会导致数据丢失或遍历异常。
修复:
添加线程锁保证安全:
# 全局变量定义处添加锁 SummarySavingList = [] SummaryLock = threading.Lock() ProblemsSavingList = [] ProblemLock = threading.Lock() # 写入列表时加锁 # 示例:在StartStreaming中添加记录时 with SummaryLock: SummarySavingList.append(f"Symbol:{Symbol},Side:LONG,Result:Takeprofit,SellTiming:{SellTiming}") # 在SummarySaver中读取处理时 while True: with SummaryLock: current_summaries = SummarySavingList.copy() for Summary in current_summaries: if not Summary in SummarySavedList: # 写入文件逻辑 with SummaryLock: SummarySavingList.remove(Summary) SummarySavedList.append(Summary)
4. 仅处理SELL方向订单
代码只针对LONG仓位的SELL方向TP/SL做了处理,完全忽略了SHORT仓位的BUY方向TP/SL(SHORT仓位的平仓操作是BUY),导致做空时反向订单无法被取消。
修复:
添加BUY方向的处理逻辑:
if Status == 'FILLED': if Side == 'SELL': # LONG仓位TP/SL处理逻辑不变 elif Side == 'BUY': # SHORT仓位TP/SL处理 if Type == 'TAKE_PROFIT_MARKET': try: Cancel = client.futures_cancel_all_open_orders(symbol=Symbol) except Exception as e: print(e) SellTiming = datetime.datetime.now().strftime("%Y-%m-%d %H-%M-%S") PositionCloser(Coin=Symbol) SummarySavingList.append(f"Symbol:{Symbol},Side:SHORT,Result:Takeprofit,SellTiming:{SellTiming}") elif Type == 'STOP_MARKET': try: Cancel = client.futures_cancel_all_open_orders(symbol=Symbol) except Exception as e: print(e) SellTiming = datetime.datetime.now().strftime("%Y-%m-%d %H-%M-%S") PositionCloser(Coin=Symbol) SummarySavingList.append(f"Symbol:{Symbol},Side:SHORT,Result:StopLoss,SellTiming:{SellTiming}")
5. PositionCloser函数逻辑缺陷
原函数只处理SELL平仓,无法平SHORT仓位;同时没有考虑币种最小交易精度,可能导致下单失败。
修复:
def PositionCloser(Coin): try: PositionInfo = client.futures_position_information(symbol=Coin)[0] position_amt = float(PositionInfo['positionAmt']) if position_amt != 0: # 获取币种最小交易精度 symbol_info = client.futures_exchange_info()['symbols'] for s in symbol_info: if s['symbol'] == Coin: min_qty = float(s['filters'][2]['minQty']) break # 调整数量到符合精度要求 quantity = abs(position_amt) quantity = round(quantity - (quantity % min_qty), 8) # 根据仓位方向选择平仓side side = 'SELL' if position_amt > 0 else 'BUY' client.futures_create_order( symbol=Coin, side=side, type='MARKET', quantity=quantity ) except Exception as e: print(f"平仓失败: {e}")
内容的提问来源于stack exchange,提问作者amj glazing
相关产品推荐
相关产品推荐

