17 KiB
17 KiB
#量化交易 #策略研究 #Python学习 #DeepSeek
这是我现在写的一个关于数字货币的一个交易脚本。请帮我完善:一是因为网络原因,如果中断执行,加入重新运行的异常处理模块。
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指标
# # 计算MACD(12/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,默认参数14,length也可以调整
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)
主要改进说明:
-
异常处理架构:
- 三层异常处理:函数级、主循环级、程序级
- 使用装饰器自动重试网络操作
- 针对不同错误类型采取不同恢复策略
-
健壮性增强:
- 指数退避 + 随机抖动重试机制
- 交易所连接重建功能
- 数据完整性检查(确保有足够数据)
- 速率限制和超时设置
-
日志系统:
- 详细记录所有操作和错误
- 同时输出到文件和终端
- 包含时间戳和日志级别
-
配置优化:
- 使用环境变量存储API密钥
- 增加API超时时间
- 启用ccxt内置速率限制
-
代码结构:
- 模块化设计,分离数据获取、处理和交易逻辑
- 函数单一职责原则
- 清晰的错误处理流程
使用建议:
-
将API密钥存储在环境变量中更安全:
export BINANCE_API_KEY='your_api_key' export BINANCE_API_SECRET='your_api_secret' -
监控日志文件
crypto_trading.log了解程序运行情况 -
对于高频交易,考虑:
- 增加
bar_num减少API调用频率 - 使用WebSocket替代轮询获取实时数据
- 添加仓位管理和风险管理模块
- 增加
这个完善版本能够在网络中断、API限制或临时错误后自动恢复,大大提高了脚本的稳定性和可靠性。