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

446 lines
17 KiB
Markdown
Raw Blame History

This file contains invisible Unicode characters

This file contains invisible Unicode characters that are indistinguishable to humans but may be processed differently by a computer. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#量化交易 #策略研究 #Python学习 #DeepSeek
---
这是我现在写的一个关于数字货币的一个交易脚本。请帮我完善:一是因为网络原因,如果中断执行,加入重新运行的异常处理模块。
![[My_strategy.py]]
```python
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错误时能够自动重试。以下是完善后的代码重点增强了健壮性和可靠性
```python
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限制或临时错误后自动恢复大大提高了脚本的稳定性和可靠性。