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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 15:19:54