Files
obsidian-vault/05 投资交易/我的第一个策略优化记录(一、异常处理).md
2026-06-23 00:24:32 +08:00

17 KiB
Raw Permalink Blame History

#量化交易 #策略研究 #Python学习 #DeepSeek


这是我现在写的一个关于数字货币的一个交易脚本。请帮我完善:一是因为网络原因,如果中断执行,加入重新运行的异常处理模块。

!My_strategy.py


import time
import pandas as pd
import ccxt
from datetime import timedelta
# import pandas_ta
import talib

# ====================================================================================================
# =====格式设置
# ====================================================================================================
pd.set_option('display.max_rows', 1000)
pd.set_option('expand_frame_repr', False)  # 当列太多时不换行
# 设置命令行输出时的列对齐功能
pd.set_option('display.unicode.ambiguous_as_wide', True)
pd.set_option('display.unicode.east_asian_width', True)

# ====================================================================================================
# =====创建ccxt交易所
# ====================================================================================================
BINANCE_CONFIG = {
                    'apiKey': '',
                    'secret': '',
                    'proxies': {'http': '127.0.0.1:7890', 'https': '127.0.0.1:7890'}
                  }
exchange = ccxt.binance(BINANCE_CONFIG)



while True:
    # list = ['DOGEUSDT', 'ETHUSDT', 'BTCUSDT']
    #
    # # for symbol in list:
    symbol = 'ETHUSDT'
    time_interval = '1m'  # 其他可以尝试的值:'1m', '5m', '15m', '30m', '1h', '2h', '1d', '1w', '1M', '1y',并不是每个交易所都支持
    bar_num = 1000  # 获取K线的数量
    params = {'symbol': symbol,  # 交易币对
              'interval': time_interval,  # 时间间隔
              'limit': bar_num}  # 数据条数

    # ====================================================================================================
    # =====获取K线数据
    # ====================================================================================================
    response = exchange.fapiPublicGetKlines(params=params)
    k_lines = pd.DataFrame(response)
    # print(k_lines)

    # =====整理K线数据
    df = pd.DataFrame(response, dtype=float)  # 将数据转换为dataframe
    df.rename(columns={0: 'MTS', 1: 'Open', 2: 'High',
                       3: 'Low', 4: 'Close', 5: 'Volume'}, inplace=True)  # 重命名
    df['candle_begin_time'] = pd.to_datetime(df['MTS'], unit='ms')  # 整理时间
    df['candle_begin_time_GMT8'] = df['candle_begin_time'] + timedelta(hours=8)  # 北京时间
    df = df[['candle_begin_time_GMT8', 'Open', 'High', 'Low', 'Close', 'Volume']]  # 整理列的顺序

    # ====================================================================================================
    # =====获取币对的最新价格
    # ====================================================================================================
    data = exchange.fapiPublicGetTickerPrice(params={'symbol': "ETHUSDT"})
    price_ETH = data['price']

    # ====================================================================================================
    # =====通过pandas-ta计算指标并加入相关指标计算列
    # ====================================================================================================
    ## 计算均线
    # df.ta.sma(length=5, append=True, col_names="SMA_5") # pandas-ta 计算SM5指标
    df['MA5'] = talib.MA(df['Close'], timeperiod=5)  # ta-lib计算MA5指标

    # # 计算14日相对强弱指数RSI
    # df.ta.rsi(length=14, append=True, col_names="RSI_14")  # pandas-ta计算RSI指标
    df['RSI_14'] = talib.RSI(df['Close'], timeperiod=14)  # ta-lib计算RSI指标

    # # 计算MACD12/26/9周期
    # df.ta.macd(fast=12, slow=26, signal=9, append=True)  # 默认列名MACD_12_26_9, MACDs_12_26_9, MACDh_12_26_9 [2,7](@ref)

    # # 计算布林带20日2倍标准差
    # df.ta.bbands(length=20, std=2, append=True)
    # pandas-ta计算布林带指标
    # 默认列名BBL_20_2.0, BBM_20_2.0, BBU_20_2.0

    ## ta-lib 计算布林带指标
    upper_band, middle_band, lower_band = talib.BBANDS(
        df['Close'],
        timeperiod=20,
        nbdevup=2,
        nbdevdn=2,
        matype=0  # SMA
    )
    df['BB_Upper'] = upper_band
    df['BB_Middle'] = middle_band
    df['BB_Lower'] = lower_band

    ## 计算ADX指标
    # df.ta.adx(length=20, append=True) # pandas-ta计算adx
    df['ADX_14'] = talib.ADX(df['High'], df['Low'], df['Close'], timeperiod=14)  # ta-lib计算ADX指标

    ## 计算ATR指标
    # df.ta.atr(length=16, append=True)     # pandas-ta计算atr默认参数14length也可以调整
    df['ATR_14'] = talib.ATR(df['High'], df['Low'], df['Close'], timeperiod=14)

    print(df)
    # df.to_csv('Ta4.csv')

    # ====================================================================================================
    # =====设置下单条件并执行
    # ====================================================================================================
    # if float(price_ETH) > float(df['Close'].iloc[-2]):
    #     print('价格上升')
    # else:
    #     print('价格下降')
    #
    # print(price_ETH)
    # print(df['Close'].iloc[-2])
    # print(df['RSI_14'].iloc[-2] > 30)

    # 以上为dataframe取数格式的测试目前使用还不太熟练

    ## 调整杠杆倍率
    # params = {'symbol': 'ETHUSDT',  # 交易币对
    #           'leverage': 10,
    #           'timestamp': int(time.time() * 1000)}
    # leverage = exchange.fapiPrivatePostLeverage(params=params)
    # print('调整开仓杠杆\n', leverage, '\n')

    ## 下单操作
    num = 0
    max_num = 3
    # 使用变量来控制下单次数

    params = {'side': 'BUY',
              'symbol': 'ETHUSDT',  # 交易币对
              'type': 'LIMIT',
              # 'price': price_ETH,
              'price': df['Close'].iloc[-1],
              # 下单价格,在限价单的时候
              # 启用将type换成LIMIT
              # 此处用变量取代了课程中的数值。以后下单的价格可能需要复杂计算再储存在变量中,做测试。
              'quantity': 0.01,  # 下单数量
              'timestamp': int(time.time() * 1000),
              'timeInForce': 'GTC'}  # https://academy.binance.com/zh/articles/understanding-the-different-order-types




    if df['RSI_14'].iloc[-2] > 30 and df['RSI_14'].iloc[-1] < 30:
        print('向下穿越rsi30执行买入')
        response = exchange.fapiPrivatePostOrder(params=params)
        order_info = pd.DataFrame(response, index=['symbol'])
        order_id = order_info['orderId'].iloc[0]
        print('下单\n', order_info, '\n')
        num += 1
        # 此后这里还可以多个条件,组合多个指标的值
        # RSI的使用作为一个参考
        # 通过现有数据(开高收低量),已经可以计算出不少指标了。

    else:
        print('不符合交易条件,监控中')
        time.sleep(30)

    if num >= max_num:
        break

## 这是我在办公室电脑上敲的注释
## 测试能否通过iCloud同步成功在家也可已经iCloud对同一个项目进行编辑。



我将为您的数字货币交易脚本添加异常处理模块确保在网络中断或API错误时能够自动重试。以下是完善后的代码重点增强了健壮性和可靠性

import time
import pandas as pd
import ccxt
from datetime import timedelta
import talib
import logging
import random
import os

# ====================================================================================================
# =====配置日志系统
# ====================================================================================================
logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s - %(levelname)s - %(message)s',
    handlers=[
        logging.FileHandler('crypto_trading.log'),
        logging.StreamHandler()
    ]
)
logger = logging.getLogger('CryptoTrader')

# ====================================================================================================
# =====格式设置
# ====================================================================================================
pd.set_option('display.max_rows', 1000)
pd.set_option('expand_frame_repr', False)
pd.set_option('display.unicode.ambiguous_as_wide', True)
pd.set_option('display.unicode.east_asian_width', True)

# ====================================================================================================
# =====创建ccxt交易所
# ====================================================================================================
def create_exchange():
    BINANCE_CONFIG = {
        'apiKey': os.getenv('BINANCE_API_KEY', ''),
        'secret': os.getenv('BINANCE_API_SECRET', ''),
        'proxies': {'http': '127.0.0.1:7890', 'https': '127.0.0.1:7890'},
        'options': {
            'defaultType': 'future',
            'adjustForTimeDifference': True
        },
        'timeout': 15000,  # 增加超时时间
        'enableRateLimit': True  # 启用速率限制
    }
    return ccxt.binance(BINANCE_CONFIG)

# ====================================================================================================
# =====重试装饰器
# ====================================================================================================
def retry(max_attempts=5, delay=2, backoff=2, exceptions=(Exception,)):
    """自动重试装饰器"""
    def decorator(func):
        def wrapper(*args, **kwargs):
            attempts = 0
            current_delay = delay
            while attempts < max_attempts:
                try:
                    return func(*args, **kwargs)
                except exceptions as e:
                    attempts += 1
                    logger.warning(f"操作失败: {str(e)}. 尝试 #{attempts}/{max_attempts}. {current_delay}秒后重试...")
                    
                    # 指数退避 + 随机抖动
                    sleep_time = current_delay + random.uniform(0, 1)
                    time.sleep(sleep_time)
                    current_delay *= backoff
            logger.error(f"操作失败 {max_attempts} 次后放弃")
            raise
        return wrapper
    return decorator

# ====================================================================================================
# =====核心函数
# ====================================================================================================
@retry(max_attempts=5, delay=3, backoff=2, exceptions=(ccxt.NetworkError, ccxt.ExchangeNotAvailable))
def fetch_klines(exchange, symbol, time_interval, bar_num):
    """获取K线数据"""
    params = {
        'symbol': symbol,
        'interval': time_interval,
        'limit': bar_num
    }
    response = exchange.fapiPublicGetKlines(params=params)
    return pd.DataFrame(response, dtype=float)

@retry(max_attempts=5, delay=2, backoff=2, exceptions=(ccxt.NetworkError, ccxt.ExchangeNotAvailable))
def fetch_price(exchange, symbol):
    """获取最新价格"""
    data = exchange.fapiPublicGetTickerPrice(params={'symbol': symbol})
    return float(data['price'])

@retry(max_attempts=3, delay=5, backoff=2, exceptions=(ccxt.NetworkError, ccxt.ExchangeNotAvailable))
def place_order(exchange, params):
    """下订单"""
    return exchange.fapiPrivatePostOrder(params=params)

# ====================================================================================================
# =====数据处理函数
# ====================================================================================================
def process_data(df):
    """处理K线数据并计算技术指标"""
    # 重命名列
    df.rename(columns={
        0: 'MTS', 1: 'Open', 2: 'High', 3: 'Low', 4: 'Close', 5: 'Volume'
    }, inplace=True)
    
    # 处理时间
    df['candle_begin_time'] = pd.to_datetime(df['MTS'], unit='ms')
    df['candle_begin_time_GMT8'] = df['candle_begin_time'] + timedelta(hours=8)
    
    # 保留必要列
    df = df[['candle_begin_time_GMT8', 'Open', 'High', 'Low', 'Close', 'Volume']]
    
    # 计算技术指标
    df['MA5'] = talib.MA(df['Close'], timeperiod=5)
    df['RSI_14'] = talib.RSI(df['Close'], timeperiod=14)
    
    # 计算布林带
    upper_band, middle_band, lower_band = talib.BBANDS(
        df['Close'], timeperiod=20, nbdevup=2, nbdevdn=2, matype=0
    )
    df['BB_Upper'] = upper_band
    df['BB_Middle'] = middle_band
    df['BB_Lower'] = lower_band
    
    # 计算ADX和ATR
    df['ADX_14'] = talib.ADX(df['High'], df['Low'], df['Close'], timeperiod=14)
    df['ATR_14'] = talib.ATR(df['High'], df['Low'], df['Close'], timeperiod=14)
    
    return df

# ====================================================================================================
# =====主交易逻辑
# ====================================================================================================
def main_trading_loop():
    """主交易循环"""
    exchange = create_exchange()
    logger.info("交易所连接成功")
    
    symbol = 'ETHUSDT'
    time_interval = '1m'
    bar_num = 1000
    max_num = 3
    num = 0
    
    while num < max_num:
        try:
            # 获取数据
            kline_df = fetch_klines(exchange, symbol, time_interval, bar_num)
            df = process_data(kline_df)
            
            # 获取最新价格
            price_ETH = fetch_price(exchange, symbol)
            logger.info(f"ETH当前价格: {price_ETH}")
            
            # 交易条件判断
            if len(df) >= 2:  # 确保有足够的数据
                rsi_prev = df['RSI_14'].iloc[-2]
                rsi_current = df['RSI_14'].iloc[-1]
                
                if rsi_prev > 30 and rsi_current < 30:
                    logger.info("向下穿越RSI30执行买入")
                    
                    # 准备订单参数
                    order_params = {
                        'side': 'BUY',
                        'symbol': symbol,
                        'type': 'LIMIT',
                        'price': df['Close'].iloc[-1],
                        'quantity': 0.01,
                        'timestamp': int(time.time() * 1000),
                        'timeInForce': 'GTC'
                    }
                    
                    # 下订单
                    response = place_order(exchange, order_params)
                    logger.info(f"订单执行成功: {response}")
                    
                    num += 1
                else:
                    logger.info("不符合交易条件,继续监控")
                    time.sleep(30)  # 常规等待
            else:
                logger.warning("数据不足,跳过本次循环")
                time.sleep(10)
                
        except ccxt.BaseError as e:
            logger.error(f"交易所API错误: {str(e)}")
            # 交易所相关错误,可能重新创建连接
            exchange = create_exchange()
            logger.info("已重新创建交易所连接")
            time.sleep(10)
            
        except Exception as e:
            logger.error(f"意外错误: {str(e)}", exc_info=True)
            time.sleep(30)
            
        except KeyboardInterrupt:
            logger.info("用户中断,程序退出")
            break

# ====================================================================================================
# =====程序入口
# ====================================================================================================
if __name__ == "__main__":
    logger.info("程序启动")
    
    while True:
        try:
            main_trading_loop()
            logger.info("交易次数达到上限,程序正常退出")
            break
        except Exception as e:
            logger.critical(f"主循环崩溃: {str(e)}", exc_info=True)
            logger.info("30秒后重启主循环...")
            time.sleep(30)

主要改进说明:

  1. 异常处理架构

    • 三层异常处理:函数级、主循环级、程序级
    • 使用装饰器自动重试网络操作
    • 针对不同错误类型采取不同恢复策略
  2. 健壮性增强

    • 指数退避 + 随机抖动重试机制
    • 交易所连接重建功能
    • 数据完整性检查(确保有足够数据)
    • 速率限制和超时设置
  3. 日志系统

    • 详细记录所有操作和错误
    • 同时输出到文件和终端
    • 包含时间戳和日志级别
  4. 配置优化

    • 使用环境变量存储API密钥
    • 增加API超时时间
    • 启用ccxt内置速率限制
  5. 代码结构

    • 模块化设计,分离数据获取、处理和交易逻辑
    • 函数单一职责原则
    • 清晰的错误处理流程

使用建议:

  1. 将API密钥存储在环境变量中更安全

    export BINANCE_API_KEY='your_api_key'
    export BINANCE_API_SECRET='your_api_secret'
    
  2. 监控日志文件crypto_trading.log了解程序运行情况

  3. 对于高频交易,考虑:

    • 增加bar_num减少API调用频率
    • 使用WebSocket替代轮询获取实时数据
    • 添加仓位管理和风险管理模块

这个完善版本能够在网络中断、API限制或临时错误后自动恢复大大提高了脚本的稳定性和可靠性。