Files
qmt_strategy/雪球6.0/源代码/交易系统_优化版.py
T
2025-11-06 10:26:02 +08:00

1377 lines
63 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
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.
# -*- coding: utf-8 -*-
"""
雪球跟单交易系统 - 优化版
严格按照业务逻辑文档规范实现
主要优化:
1. 规范化变量命名,符合业务逻辑文档术语标准
2. 重构业务流程,按照文档规范的流程组织
3. 优化数据处理逻辑,完善字段扩展和计算公式
4. 增强错误处理机制和日志输出规范
"""
import pandas as pd
import yaml
import os
import time
import datetime
import logging
from decimal import Decimal, ROUND_HALF_UP
from pyqmt_secure import pyqmt as qmt
import warnings
import subprocess
from typing import Tuple, Optional, Dict, Any
warnings.filterwarnings(action='ignore')
class XueqiuTradingSystem:
"""
雪球跟单交易系统
严格按照业务逻辑文档规范实现
"""
def __init__(self, config_file: str = None):
"""
系统初始化方法
Args:
config_file: 配置文件路径,如果为None则自动确定路径
"""
# 自动确定配置文件路径逻辑
if config_file is None:
# 获取当前脚本文件的绝对路径所在目录
current_dir = os.path.dirname(os.path.abspath(__file__))
# 配置文件位于源代码目录的上级目录,拼接完整路径
config_file = os.path.join(os.path.dirname(current_dir), '参数设置.yaml')
# 第一步:设置日志系统,为后续操作提供日志记录功能
self._setup_logging()
# 记录系统初始化开始的日志信息
self.logger.info('[系统初始化] 开始初始化雪球跟单交易系统...')
# 第二步:加载YAML配置文件,获取系统运行所需的各项参数
self.config = self._load_config(config_file)
# 第三步:根据配置文件初始化系统核心参数(如跟单比例、交易时间等)
self._init_system_parameters()
# 第四步:初始化QMT量化交易接口,建立与券商交易系统的连接
self._init_qmt_interface()
# 第五步:初始化数据管理模块,设置数据存储和处理相关参数
self._init_data_management()
# 记录系统初始化完成的日志信息
self.logger.info('[系统初始化] 雪球跟单交易系统初始化完成')
def _setup_logging(self):
"""
设置日志系统
按照业务逻辑文档规范设置日志级别和输出格式
"""
# 创建日志目录
log_dir = os.path.dirname(os.path.abspath(__file__))
log_file = os.path.join(log_dir, 'trading_system.log')
# 设置日志格式
log_format = '%(asctime)s - %(levelname)s - %(message)s'
date_format = '%Y-%m-%d %H:%M:%S'
# 配置日志处理器
handlers = [
logging.StreamHandler(), # 控制台输出
logging.FileHandler(log_file, encoding='utf-8') # 文件输出
]
# 设置日志级别
log_level = self.config.get('日志级别', 'INFO') if hasattr(self, 'config') else 'INFO'
level = getattr(logging, log_level.upper(), logging.INFO)
logging.basicConfig(
level=level,
format=log_format,
datefmt=date_format,
handlers=handlers,
force=True # 强制重新配置
)
self.logger = logging.getLogger(__name__)
self.logger.info(f'[日志系统] 日志级别设置为: {log_level}')
self.logger.info(f'[日志系统] 日志文件路径: {log_file}')
def _load_config(self, config_file: str) -> Dict[str, Any]:
"""
加载配置文件
Args:
config_file: 配置文件路径
Returns:
配置字典
"""
try:
if os.path.exists(config_file):
with open(config_file, 'r', encoding='utf-8') as f:
config = yaml.safe_load(f)
self.logger.info(f'[配置加载] 成功加载配置文件: {config_file}')
return config
else:
self.logger.warning(f'[配置加载] 配置文件不存在: {config_file},使用默认配置')
return {}
except Exception as e:
self.logger.error(f'[配置加载] 读取配置文件失败: {str(e)}')
return {}
def _init_system_parameters(self):
"""
初始化系统参数
"""
# 从配置文件读取参数
self.account_follow_ratio = self.config.get('账户跟单比例(%)', 0.0)
self.portfolio_configs = self.config.get('组合配置', [])
# 确保组合配置中的跟单比例字段名称一致
for config in self.portfolio_configs:
if '组合跟单比例(%)' not in config and '跟单比例' in config:
config['组合跟单比例(%)'] = config['跟单比例'] * 100
# 参数值范围验证
if not (0 <= self.account_follow_ratio <= 100):
self.logger.warning(f'[参数验证] 账户跟单比例(%)值 {self.account_follow_ratio} 超出范围0-100,将使用默认值0')
self.account_follow_ratio = 0.0
for config in self.portfolio_configs:
ratio = config.get('组合跟单比例(%)', 0.0)
if not (0 <= ratio <= 100):
self.logger.warning(f'[参数验证] 组合 {config.get("组合名字", "未知")} 的跟单比例(%)值 {ratio} 超出范围0-100,将使用默认值0')
config['组合跟单比例(%)'] = 0
# 验证所有组合跟单比例总和不超过100%
total_ratio = sum(config.get('组合跟单比例(%)', 0.0) for config in self.portfolio_configs)
if total_ratio > 100:
self.logger.warning(f'[参数验证] 所有组合跟单比例总和 {total_ratio}% 超过100%,将按比例缩放到100%以内')
# 按比例缩放所有组合的跟单比例
scale_factor = 100.0 / total_ratio
for config in self.portfolio_configs:
original_ratio = config.get('组合跟单比例(%)', 0.0)
scaled_ratio = round(original_ratio * scale_factor, 2)
config['组合跟单比例(%)'] = scaled_ratio
self.logger.info(f'[参数调整] 组合 {config.get("组合名字", "未知")} 跟单比例从 {original_ratio}% 调整为 {scaled_ratio}%')
# 重新计算总和并记录
new_total = sum(config.get('组合跟单比例(%)', 0.0) for config in self.portfolio_configs)
self.logger.info(f'[参数调整] 调整后所有组合跟单比例总和: {new_total}%')
else:
self.logger.info(f'[参数验证] 所有组合跟单比例总和: {total_ratio}%,符合要求')
self.loop_interval = 3 # 循环间隔(秒)
# 数据文件路径配置 - 使用绝对路径
current_dir = os.path.dirname(os.path.abspath(__file__))
project_root = os.path.dirname(current_dir)
auxiliary_dir = os.path.join(project_root, '辅助文件')
self.data_file_path = os.path.join(auxiliary_dir, '数据.csv') # 最终处理后的数据文件
self.raw_data_file_path = os.path.join(auxiliary_dir, '雪球原始数据.csv') # 雪球原始数据文件
# 已处理记录ID集合
self.processed_record_ids = set()
# 统计信息
self.execution_count = 0
self.total_new_records = 0
self.logger.info(f'[系统参数] 账户跟单比例(%): {self.account_follow_ratio}%')
self.logger.info(f'[系统参数] 组合配置数量: {len(self.portfolio_configs)}')
def _init_qmt_interface(self):
"""
初始化QMT交易接口
"""
try:
qmt_path = self.config.get('QMT路径', 'C:\\国金QMT交易端模拟\\userdata_mini')
account_id = self.config.get('资金账号', '39972702')
self.logger.info(f'[QMT初始化] QMT路径: {qmt_path}')
self.logger.info(f'[QMT初始化] 资金账号: {account_id}')
self.qmt = qmt(path=qmt_path, acc=account_id)
self.qmt_available = True
self.logger.info('[QMT初始化] QMT交易接口初始化成功')
except Exception as e:
self.logger.warning(f'[QMT初始化] QMT交易接口初始化失败: {str(e)}')
self.logger.warning('[QMT初始化] QMT连接失败不影响数据处理和计算逻辑,系统将继续运行(仅交易功能不可用)')
self.qmt = None
self.qmt_available = False
def _init_data_management(self):
"""
初始化数据管理
"""
try:
if os.path.exists(self.data_file_path):
df = pd.read_csv(self.data_file_path, encoding='utf-8-sig')
# 检查文件是否只有表头(空数据)
if len(df) == 0:
self.logger.info('[数据管理] 数据文件为空,重置已处理记录ID集合')
self.processed_record_ids = set()
elif '记录ID' in df.columns:
self.processed_record_ids = set(df['记录ID'].astype(str))
self.logger.info(f'[数据管理] 已加载 {len(self.processed_record_ids)} 个已处理记录ID')
else:
self.logger.warning('[数据管理] 数据文件中未找到记录ID列')
self.processed_record_ids = set()
else:
self.logger.info('[数据管理] 数据文件不存在,将从空数据开始')
self.processed_record_ids = set()
except Exception as e:
self.logger.error(f'[数据管理] 初始化数据管理失败: {str(e)}')
self.processed_record_ids = set()
def is_trading_time(self) -> Tuple[bool, str]:
"""
判断当前是否为交易时间
Returns:
(是否为交易时间, 时间状态信息)
"""
now = datetime.datetime.now()
current_time = now.time()
# 交易时间:周一至周五 9:15-15:00(测试模式)
trading_start = datetime.time(9, 15)
trading_end = datetime.time(15, 00)
# 检查是否为工作日
if now.weekday() >= 5: # 周六=5, 周日=6
return False, "当前为周末,非交易时间"
# 检查时间范围
if trading_start <= current_time <= trading_end:
return True, "当前为交易时间"
else:
return False, f"当前时间 {current_time.strftime('%H:%M:%S')} 不在交易时间范围内(00:00-23:59"
def execute_xueqiu_script(self) -> Tuple[bool, Optional[str]]:
"""
执行雪球数据获取脚本
Returns:
(是否成功, 错误信息)
"""
try:
current_dir = os.path.dirname(os.path.abspath(__file__))
xueqiu_path = os.path.join(current_dir, 'xueqiu.py')
project_root = os.path.dirname(current_dir)
result = subprocess.run(
['python', xueqiu_path],
capture_output=True,
text=True,
encoding='gbk',
cwd=project_root
)
if result.returncode == 0:
self.logger.info('[雪球数据] 雪球数据获取成功')
return True, None
else:
self.logger.error(f'[雪球数据] 雪球数据获取失败: {result.stderr}')
return False, result.stderr
except Exception as e:
self.logger.error(f'[雪球数据] 执行雪球脚本异常: {str(e)}')
return False, str(e)
def get_account_total_assets(self) -> Optional[float]:
"""
获取账户总资产
Returns:
账户总资产,获取失败返回None
"""
try:
if not self.qmt_available or self.qmt is None:
self.logger.warning('[账户资产] QMT接口不可用,返回None')
return None
asset_df = self.qmt.query_stock_asset()
if asset_df is None or asset_df.empty:
self.logger.warning('[账户资产] 查询账户资产返回空数据,返回None')
return None
total_assets = asset_df['总资产'].iloc[0]
self.logger.info(f'[账户资产] 账户总资产: {total_assets:.2f}')
return total_assets
except Exception as e:
self.logger.warning(f'[账户资产] 获取账户总资产失败: {str(e)},返回None')
return None
def convert_stock_code_format(self, original_code: str) -> str:
"""
转换股票代码格式
从 SH600109 格式转换为 600109.SH 格式
Args:
original_code: 原始股票代码
Returns:
转换后的股票代码
"""
if original_code.startswith('SH'):
return original_code[2:] + '.SH'
elif original_code.startswith('SZ'):
return original_code[2:] + '.SZ'
else:
return original_code
def calculate_adjustment_ratio(self, target_weight: float, previous_weight: float) -> float:
"""
计算调仓幅度
调仓幅度 = 目标权重 - 前期调整权重
Args:
target_weight: 目标权重
previous_weight: 前期调整权重
Returns:
调仓幅度
"""
# 空值处理:当目标权重或前期调整权重为空时,视为0
if pd.isna(target_weight):
target_weight = 0.0
if pd.isna(previous_weight):
previous_weight = 0.0
adjustment_ratio = target_weight - previous_weight
return adjustment_ratio
def determine_trade_direction(self, adjustment_ratio: float) -> str:
"""
确定交易方向
Args:
adjustment_ratio: 调仓幅度
Returns:
交易方向:买入/卖出/无操作
"""
if adjustment_ratio > 0:
return "买入"
elif adjustment_ratio < 0:
return "卖出"
else:
return "无操作"
def calculate_follow_shares(self, account_total_assets: float, account_follow_ratio: float,
portfolio_ratio: float, adjustment_ratio: float,
latest_price: float, stock_code: str = "") -> int:
"""
计算跟单股数
Args:
account_total_assets: 账户总资产
account_follow_ratio: 账户跟单比例(%)
portfolio_ratio: 组合跟单比例(%)
adjustment_ratio: 调仓幅度(%)
latest_price: 标的最新价
stock_code: 股票代码,用于判断是否为可转债
Returns:
跟单股数(可转债向下取整到10股的整数倍,普通股票向下取整到100股的整数倍)
"""
# 跟单金额 = 账户总资产 × 账户跟随比例/100 × 组合跟随比例/100 × 调仓幅度的绝对值/100
follow_amount = account_total_assets * (account_follow_ratio / 100) * (portfolio_ratio / 100) * (abs(adjustment_ratio) / 100)
# 判断是否为可转债(11或12开头)
is_convertible_bond = stock_code.startswith('11') or stock_code.startswith('12')
if is_convertible_bond:
# 可转债:跟单股数 = 跟单金额 // (最新价 × 10) × 10
follow_shares = int(follow_amount // (latest_price * 10)) * 10
else:
# 普通股票:跟单股数 = 跟单金额 // (最新价 × 100) × 100
follow_shares = int(follow_amount // (latest_price * 100)) * 100
return max(0, follow_shares) # 确保不为负数
def get_portfolio_name(self, portfolio_id: str) -> str:
"""
根据组合ID获取组合名字
Args:
portfolio_id: 组合ID
Returns:
组合名字
"""
for config in self.portfolio_configs:
if config.get('组合ID') == portfolio_id:
return config.get('组合名字', f'组合{portfolio_id}')
return f'组合{portfolio_id}'
def get_portfolio_ratio(self, portfolio_id: str) -> float:
"""
根据组合ID获取组合跟单比例(%)
Args:
portfolio_id: 组合ID
Returns:
组合跟单比例(%)
"""
for config in self.portfolio_configs:
if config.get('组合ID') == portfolio_id:
return config.get('组合跟单比例(%)', 0.0)
return 0.0
def extend_data_fields(self, df: pd.DataFrame) -> pd.DataFrame:
"""
扩展数据字段
按照业务逻辑文档规范添加必要字段
Args:
df: 原始数据DataFrame
Returns:
扩展后的DataFrame
"""
import numpy as np
# 删除旧版本的跟单比例(%)列(已被组合跟单比例(%)替代)
if '跟单比例(%)' in df.columns:
df = df.drop('跟单比例(%)', axis=1)
self.logger.info('[字段扩展] 已删除旧版本跟单比例(%)列')
# 3.1 配置参数字段
required_fields = [
'账户跟单比例(%)', '组合名字', '组合ID', '组合跟单比例(%)',
'账户总资产', '当前权重', '调仓幅度', '交易方向', '代码',
'跟单股数', '是否委托', '调仓时间'
]
for field in required_fields:
if field not in df.columns:
if field in ['组合名字', '交易方向', '代码', '调仓时间']:
df[field] = ''
elif field == '是否委托':
df[field] = '' # 初始化为"否"
elif field == '账户总资产':
df[field] = np.nan # 账户总资产初始化为NAN
else:
df[field] = 0.0
self.logger.info(f'[字段扩展] 已添加字段: {field}')
# 添加调仓时间字段:将更新时间的毫秒时间戳转换为YYYYMMDDHHMMSS格式
if '更新时间' in df.columns:
try:
# 将毫秒时间戳转换为datetime,然后格式化为YYYYMMDDHHMMSS
df['调仓时间'] = pd.to_datetime(df['更新时间'], unit='ms').dt.strftime('%Y%m%d%H%M%S')
self.logger.info('[字段扩展] 已添加调仓时间字段,格式:YYYYMMDDHHMMSS')
except Exception as e:
self.logger.error(f'[字段扩展] 转换调仓时间失败: {str(e)}')
df['调仓时间'] = ''
# 数据计算完成,记录总体统计信息
self.logger.info(f'[数据计算] 数据计算完成,共处理{len(df)}条记录')
return df
def process_data_calculations(self, df: pd.DataFrame) -> pd.DataFrame:
"""
处理数据计算
按照业务逻辑文档规范进行各项计算
Args:
df: 数据DataFrame
Returns:
处理后的DataFrame
"""
import numpy as np
# 预先获取账户总资产,避免每行都重复获取
account_total_assets = np.nan # 默认为NAN
if self.qmt is not None:
try:
assets = self.get_account_total_assets()
if assets and assets > 0:
account_total_assets = assets
self.logger.info(f'[数据计算] 获取账户总资产: {account_total_assets}')
else:
self.logger.info('[数据计算] QMT连接正常但获取资产失败,使用NAN')
except Exception:
self.logger.info('[数据计算] 获取账户总资产异常,使用NAN')
else:
self.logger.info('[数据计算] QMT未连接,使用NAN作为账户总资产')
# 批量处理所有数据,避免逐行调用QMT接口
for index, row in df.iterrows():
try:
# 2.2 调仓幅度计算 - 不依赖QMT
target_weight = row.get('目标权重', 0.0)
previous_weight = row.get('前期调整权重', 0.0)
adjustment_ratio = self.calculate_adjustment_ratio(target_weight, previous_weight)
# 调仓幅度保留2位小数并四舍五入
df.loc[index, '调仓幅度'] = round(adjustment_ratio, 2)
# 2.3 交易方向确定 - 不依赖QMT
trade_direction = self.determine_trade_direction(adjustment_ratio)
df.loc[index, '交易方向'] = trade_direction
# 3.3 股票代码格式转换 - 不依赖QMT
original_code = row.get('股票代码', '')
if original_code:
converted_code = self.convert_stock_code_format(original_code)
df.loc[index, '代码'] = converted_code
# 账户跟单比例 - 不依赖QMT,直接从配置获取
if hasattr(self, 'account_follow_ratio') and self.account_follow_ratio is not None and self.account_follow_ratio > 0:
df.loc[index, '账户跟单比例(%)'] = self.account_follow_ratio
else:
self.logger.warning('[数据计算] 账户跟单比例未正确配置,使用默认值10%')
df.loc[index, '账户跟单比例(%)'] = 10.0
# 组合相关信息 - 不依赖QMT,直接从配置获取
portfolio_id = row.get('组合ID', '')
if portfolio_id:
portfolio_name = self.get_portfolio_name(portfolio_id)
portfolio_ratio = self.get_portfolio_ratio(portfolio_id)
df.loc[index, '组合名字'] = portfolio_name
# 组合跟单比例 - 不依赖QMT,修复逻辑确保正确获取配置值
if portfolio_ratio is not None and portfolio_ratio > 0:
df.loc[index, '组合跟单比例(%)'] = portfolio_ratio
# 减少日志输出,只在第一次记录时输出
if not hasattr(self, '_logged_portfolio_ratios'):
self._logged_portfolio_ratios = set()
if portfolio_id not in self._logged_portfolio_ratios:
self.logger.info(f'[数据计算] 组合ID {portfolio_id} 跟单比例: {portfolio_ratio}%')
self._logged_portfolio_ratios.add(portfolio_id)
else:
# 如果配置中没有找到,使用默认值但记录警告
if not hasattr(self, '_logged_portfolio_warnings'):
self._logged_portfolio_warnings = set()
if portfolio_id not in self._logged_portfolio_warnings:
self.logger.warning(f'[数据计算] 组合ID {portfolio_id} 在配置中未找到,使用默认值50%')
self._logged_portfolio_warnings.add(portfolio_id)
df.loc[index, '组合跟单比例(%)'] = 50.0
else:
df.loc[index, '组合名字'] = ''
df.loc[index, '组合跟单比例(%)'] = 50.0 # 默认值
# 是否委托 - 不依赖QMT,初始化为"否"
if '是否委托' not in df.columns or pd.isna(df.loc[index, '是否委托']):
df.loc[index, '是否委托'] = ''
# 3.2 账户资产字段 - 使用预先获取的值,不再重复调用QMT
df.loc[index, '账户总资产'] = account_total_assets
# 4.1 跟单股数计算 - 现在所有关键字段都有值,可以正常计算
current_account_assets = account_total_assets # 使用预先获取的值
current_account_ratio = df.loc[index, '账户跟单比例(%)']
current_portfolio_ratio = df.loc[index, '组合跟单比例(%)']
# 检查关键参数是否为NAN,如果任一为NAN则跟单股数也设为NAN
if (pd.isna(current_account_assets) or
pd.isna(current_account_ratio) or
pd.isna(current_portfolio_ratio)):
df.loc[index, '跟单股数'] = np.nan
continue
# 计算跟单股数 - 不再过滤调仓幅度为0的记录,处理所有有股票代码的记录
if df.loc[index, '代码']:
try:
# 获取股票真实价格
stock_code = df.loc[index, '代码']
try:
# 使用QMT接口获取最新价
latest_price = self.qmt.get_last_price(stock_code)
if latest_price is None or latest_price <= 0:
raise ValueError(f"获取到无效价格: {latest_price}")
self.logger.info(f'[价格获取] {stock_code} 最新价格: {latest_price}')
except Exception as price_error:
# 价格获取失败时的处理
self.logger.warning(f'[价格获取] 获取 {stock_code} 最新价失败: {str(price_error)},跳过该股票下单')
df.loc[index, '跟单股数'] = 0
df.loc[index, '是否委托'] = ''
continue
# 使用真实价格计算跟单股数
follow_shares = self.calculate_follow_shares(
current_account_assets, current_account_ratio,
current_portfolio_ratio, adjustment_ratio, latest_price, stock_code
)
df.loc[index, '跟单股数'] = follow_shares
except Exception as e:
self.logger.error(f'[数据计算] 计算跟单股数失败: {str(e)}')
df.loc[index, '跟单股数'] = 0
else:
# 无股票代码时,跟单股数为0,但仍保留记录(减少日志输出)
df.loc[index, '跟单股数'] = 0
except Exception as e:
self.logger.error(f'[数据计算] 处理第{index}行数据时发生异常: {str(e)}')
# 确保即使出错也有基本的字段值
df.loc[index, '调仓幅度'] = 0.0
df.loc[index, '交易方向'] = '无操作'
df.loc[index, '跟单股数'] = 0
df.loc[index, '是否委托'] = ''
continue
return df
def execute_trade_order(self, symbol: str, volume: int, trade_direction: str,
strategy_name: str) -> Optional[int]:
"""
执行交易订单
Args:
symbol: 股票代码
volume: 交易数量
trade_direction: 交易方向
strategy_name: 策略名称
Returns:
订单ID,失败返回None
"""
try:
if not self.qmt_available or self.qmt is None:
self.logger.warning(f'[交易执行] QMT接口不可用,跳过交易: {trade_direction} {symbol} {volume}')
return None
if trade_direction == "买入":
order_id = self.qmt.buy(symbol, volume, strategy_name=strategy_name)
elif trade_direction == "卖出":
order_id = self.qmt.sell(symbol, volume, strategy_name=strategy_name)
else:
self.logger.info(f'[交易执行] {symbol} 无需操作')
return None
if order_id and order_id > 0:
self.logger.info(f'[交易执行] {trade_direction} {symbol} {volume}股 成功,订单ID: {order_id}')
return order_id
else:
self.logger.error(f'[交易执行] {trade_direction} {symbol} {volume}股 失败')
return None
except Exception as e:
self.logger.error(f'[交易执行] 执行交易异常: {str(e)}')
return None
def update_trade_status(self, df: pd.DataFrame, index: int, success: bool) -> pd.DataFrame:
"""
更新交易状态
Args:
df: 数据DataFrame
index: 记录索引
success: 是否成功
Returns:
更新后的DataFrame
"""
if success:
df.loc[index, '是否委托'] = ''
self.logger.info(f'[状态更新] 第{index+1}条记录交易状态已更新为:是')
else:
df.loc[index, '是否委托'] = ''
self.logger.warning(f'[状态更新] 第{index+1}条记录交易失败,状态保持为:否')
return df
def save_data_with_backup(self, df: pd.DataFrame) -> bool:
"""
保存数据(已取消备份功能)
Args:
df: 要保存的DataFrame
Returns:
保存是否成功
"""
try:
# 创建副本以避免修改原始数据
df_to_save = df.copy()
# 将所有类型的空值(NaN、None、空字符串)替换为"NAN"字符串
import numpy as np
df_to_save = df_to_save.replace([np.nan, None, ''], 'NAN')
# 直接保存数据,不再创建备份文件
df_to_save.to_csv(self.data_file_path, index=False, encoding='utf-8-sig')
self.logger.info(f'[数据保存] 数据保存成功,共保存{len(df)}条数据到文件: {self.data_file_path}')
return True
except Exception as e:
self.logger.error(f'[数据保存] 保存数据失败: {str(e)}')
return False
def filter_duplicate_records(self, df: pd.DataFrame) -> pd.DataFrame:
"""
基于记录ID和调仓ID进行智能去重处理
修复重复下单问题:保持已委托记录的状态,只对真正的新记录设置""状态
Args:
df: 原始数据DataFrame(通常来自雪球原始数据)
Returns:
去重后需要处理的记录DataFrame
"""
if df.empty or '记录ID' not in df.columns:
self.logger.warning('[数据去重] 数据为空或缺少记录ID列')
return pd.DataFrame()
# 读取现有的已处理数据文件,获取当前的委托状态
existing_df = pd.DataFrame()
if os.path.exists(self.data_file_path):
try:
existing_df = pd.read_csv(self.data_file_path, encoding='utf-8-sig')
self.logger.info(f'[数据去重] 成功读取现有数据文件,共{len(existing_df)}条记录')
except Exception as e:
self.logger.error(f'[数据去重] 读取现有数据文件失败: {str(e)}')
existing_df = pd.DataFrame()
# 转换记录ID为字符串进行比较
df['记录ID_str'] = df['记录ID'].astype(str)
# 创建唯一标识:股票代码 + 调仓ID,用于更精确的重复判断
if '调仓ID' in df.columns and '股票代码' in df.columns:
df['唯一标识'] = df['股票代码'].astype(str) + '_' + df['调仓ID'].astype(str)
self.logger.info('[数据去重] 使用股票代码+调仓ID作为唯一标识')
else:
# 如果没有调仓ID,则使用记录ID作为唯一标识
df['唯一标识'] = df['记录ID'].astype(str)
self.logger.warning('[数据去重] 缺少调仓ID字段,使用记录ID作为唯一标识')
# 筛选出未处理过的记录ID(基于内存中的processed_record_ids
new_records = df[~df['记录ID_str'].isin(self.processed_record_ids)].copy()
if len(new_records) > 0:
# 删除临时列
new_records = new_records.drop('记录ID_str', axis=1)
# 智能设置委托状态:检查是否在现有数据中已经存在且已委托
if '是否委托' not in new_records.columns:
# 如果原始数据中没有委托状态字段,则添加该字段
new_records['是否委托'] = ''
self.logger.info('[数据去重] 原始数据无委托状态字段,为所有新记录初始化为""')
else:
# 如果原始数据中有委托状态字段,先填充空值为"否"
new_records['是否委托'] = new_records['是否委托'].fillna('')
# 关键修复:检查现有数据中是否已经存在相同的记录且已委托
if not existing_df.empty and '唯一标识' in existing_df.columns and '是否委托' in existing_df.columns:
# 为现有数据也创建唯一标识
if '调仓ID' in existing_df.columns and '股票代码' in existing_df.columns:
existing_df['唯一标识'] = existing_df['股票代码'].astype(str) + '_' + existing_df['调仓ID'].astype(str)
else:
existing_df['唯一标识'] = existing_df['记录ID'].astype(str)
# 检查每条新记录是否在现有数据中已经存在且已委托
for idx, row in new_records.iterrows():
unique_id = row['唯一标识']
# 在现有数据中查找相同唯一标识的记录
existing_record = existing_df[existing_df['唯一标识'] == unique_id]
if not existing_record.empty:
# 如果找到相同记录,检查其委托状态
existing_status = existing_record.iloc[0]['是否委托']
if existing_status == '':
# 如果已经委托,保持"是"状态,避免重复下单
new_records.loc[idx, '是否委托'] = ''
self.logger.info(f'[数据去重] 记录{unique_id}已存在且已委托,保持""状态')
else:
# 如果未委托,保持"否"状态
new_records.loc[idx, '是否委托'] = ''
self.logger.info(f'[数据去重] 记录{unique_id}已存在但未委托,保持""状态')
else:
# 如果是全新记录,设置为"否"
new_records.loc[idx, '是否委托'] = ''
self.logger.info(f'[数据去重] 记录{unique_id}为全新记录,设置为""状态')
# 删除临时的唯一标识列
if '唯一标识' in new_records.columns:
new_records = new_records.drop('唯一标识', axis=1)
# 更新已处理记录ID集合
new_record_ids = set(new_records['记录ID'].astype(str))
self.processed_record_ids.update(new_record_ids)
# 统计委托状态分布
status_counts = new_records['是否委托'].value_counts()
self.logger.info(f'[数据去重] 发现{len(new_records)}条记录,委托状态分布: {dict(status_counts)}')
else:
self.logger.info('[数据去重] 没有发现新记录')
return new_records
def query_and_display_trades(self):
"""
查询并显示成交信息
按照业务逻辑文档规范输出交易结果
"""
try:
if not self.qmt_available or self.qmt is None:
self.logger.warning('[成交查询] QMT接口不可用,跳过成交信息查询')
return
# 查询今日成交记录
trades_df = self.qmt.query_stock_trades()
if trades_df is None or trades_df.empty:
self.logger.info('[成交查询] 当前没有成交记录')
return
# 处理成交时间格式
if '成交时间' in trades_df.columns:
trades_df['成交时间'] = pd.to_datetime(trades_df['成交时间'], unit='s', utc=True)
trades_df['成交时间'] = trades_df['成交时间'].dt.tz_convert('Asia/Shanghai')
trades_df['成交时间'] = trades_df['成交时间'].dt.strftime("%Y-%m-%d %H:%M:%S")
# 按策略名称(组合)分组显示
if '策略名称' in trades_df.columns:
for strategy_name, group in trades_df.groupby('策略名称'):
if strategy_name and strategy_name != '未知':
self.logger.info(f'\n--- {strategy_name} 成交信息 ---')
for index, trade in group.iterrows():
order_type = '买入' if trade.get('委托类型') == 23 else '卖出' if trade.get('委托类型') == 24 else '未知'
self.logger.info(
f'[成交] {trade.get("证券代码", "未知")} | {order_type} | '
f'数量:{trade.get("成交数量", 0)} | 价格:{trade.get("成交均价", 0):.2f} | '
f'金额:{trade.get("成交金额", 0):.2f} | 时间:{trade.get("成交时间", "未知")}'
)
self.logger.info(f'[成交统计] 共 {len(trades_df)} 条成交记录')
except Exception as e:
self.logger.error(f'[成交查询] 查询成交信息失败: {str(e)}')
def handle_critical_error(self, operation: str, error: Exception, retry_count: int = 0) -> bool:
"""
处理关键错误
Args:
operation: 操作名称
error: 异常对象
retry_count: 重试次数
Returns:
是否应该继续运行
"""
error_msg = str(error)
# 记录错误详情
self.logger.error(f'[关键错误] {operation} 发生错误: {error_msg}')
self.logger.error(f'[关键错误] 错误类型: {type(error).__name__}')
# 根据错误类型决定处理策略
import requests
if isinstance(error, (ConnectionError, requests.exceptions.RequestException)):
self.logger.warning(f'[网络错误] {operation} 网络连接问题,重试次数: {retry_count}')
return retry_count < 3 # 网络错误最多重试3次
elif isinstance(error, FileNotFoundError):
self.logger.error(f'[文件错误] {operation} 文件不存在: {error_msg}')
return False # 文件错误不重试
elif isinstance(error, PermissionError):
self.logger.error(f'[权限错误] {operation} 权限不足: {error_msg}')
return False # 权限错误不重试
elif isinstance(error, (ValueError, TypeError)):
self.logger.error(f'[数据错误] {operation} 数据格式错误: {error_msg}')
return True # 数据错误可以继续运行
else:
self.logger.error(f'[未知错误] {operation} 未知错误类型: {error_msg}')
return retry_count < 2 # 未知错误最多重试2次
def handle_network_exception(self, operation: str, max_retries: int = 3) -> bool:
"""
处理网络异常的重试机制
Args:
operation: 操作名称
max_retries: 最大重试次数
Returns:
是否成功处理
"""
for retry in range(max_retries):
try:
self.logger.warning(f'[网络重试] {operation}{retry + 1}次重试')
time.sleep(2 ** retry) # 指数退避
return True
except Exception as e:
if retry == max_retries - 1:
self.logger.error(f'[网络异常] {operation} 重试{max_retries}次后仍然失败: {str(e)}')
return False
continue
return False
def validate_system_health(self) -> bool:
"""
系统健康检查
Returns:
系统是否健康
"""
try:
self.logger.info('[系统检查] 开始系统健康检查')
# 检查配置文件
if not hasattr(self, 'config') or not self.config:
self.logger.error('[系统检查] 配置文件未加载')
return False
# 检查必要的配置项
required_configs = ['cookie列表', 'QMT路径', '资金账号', '组合配置']
for config_key in required_configs:
if config_key not in self.config:
self.logger.error(f'[系统检查] 缺少必要配置项: {config_key}')
return False
# 检查数据文件路径
data_path = self.config.get('数据保存路径', '辅助文件\\数据.csv')
data_dir = os.path.dirname(data_path)
if not os.path.exists(data_dir):
self.logger.warning(f'[系统检查] 数据目录不存在,将创建: {data_dir}')
os.makedirs(data_dir, exist_ok=True)
# 检查QMT路径
qmt_path = self.config.get('QMT路径')
if qmt_path and not os.path.exists(qmt_path):
self.logger.warning(f'[系统检查] QMT路径不存在: {qmt_path}')
# 检查组合配置
portfolio_configs = self.config.get('组合配置', [])
if not portfolio_configs:
self.logger.error('[系统检查] 未配置任何投资组合')
return False
for i, portfolio in enumerate(portfolio_configs):
if '组合名字' not in portfolio or '组合ID' not in portfolio:
self.logger.error(f'[系统检查] 组合配置{i+1}缺少必要字段')
return False
self.logger.info('[系统检查] 系统健康检查通过')
return True
except Exception as e:
self.logger.error(f'[系统检查] 健康检查失败: {str(e)}')
return False
def validate_data_integrity(self, df: pd.DataFrame) -> Tuple[bool, str]:
"""
验证数据完整性
Args:
df: 数据DataFrame
Returns:
(是否有效, 错误信息)
"""
required_columns = ['记录ID', '股票代码', '股票名称', '目标权重', '组合ID']
# 检查必需列
missing_columns = [col for col in required_columns if col not in df.columns]
if missing_columns:
return False, f"缺少必需列: {', '.join(missing_columns)}"
# 检查数据类型
for index, row in df.iterrows():
try:
# 检查记录ID
if pd.isna(row['记录ID']) or str(row['记录ID']).strip() == '':
return False, f"{index+1}行记录ID为空"
# 检查目标权重
if pd.isna(row['目标权重']):
return False, f"{index+1}行目标权重为空"
target_weight = float(row['目标权重'])
if target_weight < 0 or target_weight > 100:
self.logger.warning(f'[数据验证] 第{index+1}行目标权重异常: {target_weight}')
except (ValueError, TypeError) as e:
return False, f"{index+1}行数据类型错误: {str(e)}"
return True, "数据验证通过"
def process_trading_data(self) -> bool:
"""
处理交易数据的主要业务流程
严格按照业务逻辑文档规范执行,增强错误处理
Returns:
处理是否成功
"""
try:
# 1. 数据获取与处理 - 2.1 雪球信号获取
self.logger.info('[数据获取] 开始获取雪球数据...')
# 调用雪球数据获取脚本
success, error = self.execute_xueqiu_script()
if not success:
self.logger.error(f'[数据获取] 雪球数据获取失败: {error}')
# 尝试网络重试
if not self.handle_network_exception("xueqiu_data"):
self.logger.error('[数据获取] 雪球数据获取重试失败,使用现有数据继续处理')
else:
self.logger.info('[数据获取] 雪球数据获取成功')
if not os.path.exists(self.data_file_path):
self.logger.warning('[数据处理] 数据文件不存在,雪球数据获取可能失败,使用空数据继续处理')
# 创建空的数据文件以确保系统正常运行
empty_df = pd.DataFrame(columns=[
'记录ID', '目标权重', '前期调整权重', '调仓幅度', '交易方向',
'股票代码', '代码', '跟单股数', '是否委托', '账户总资产',
'账户跟单比例(%)', '组合名字', '组合ID', '组合跟单比例(%)'
])
try:
empty_df.to_csv(self.data_file_path, index=False, encoding='utf-8-sig')
self.logger.info('[数据处理] 已创建空数据文件,系统继续运行')
except Exception as e:
self.logger.error(f'[数据处理] 创建空数据文件失败: {str(e)}')
return False
# 智能数据读取与状态合并逻辑
try:
df = pd.DataFrame() # 初始化空DataFrame
# 1. 优先读取雪球原始数据文件(最新的信号数据)
if os.path.exists(self.raw_data_file_path):
raw_df = pd.read_csv(self.raw_data_file_path, encoding='utf-8-sig')
self.logger.info(f'[数据读取] 成功读取雪球原始数据文件,共{len(raw_df)}条记录')
# 2. 读取已处理数据文件(包含委托状态信息)
processed_df = pd.DataFrame()
if os.path.exists(self.data_file_path):
try:
processed_df = pd.read_csv(self.data_file_path, encoding='utf-8-sig')
self.logger.info(f'[数据读取] 成功读取已处理数据文件,共{len(processed_df)}条记录')
except Exception as e:
self.logger.warning(f'[数据读取] 读取已处理数据文件失败: {str(e)}')
processed_df = pd.DataFrame()
# 3. 智能合并原始数据与已处理数据的状态
if not processed_df.empty:
# 为原始数据添加默认的委托状态字段
if '是否委托' not in raw_df.columns:
raw_df['是否委托'] = '' # 默认为未委托
self.logger.info('[数据合并] 为原始数据添加默认委托状态字段')
# 创建唯一标识用于状态匹配
if '调仓ID' in raw_df.columns and '股票代码' in raw_df.columns:
raw_df['唯一标识'] = raw_df['股票代码'].astype(str) + '_' + raw_df['调仓ID'].astype(str)
else:
raw_df['唯一标识'] = raw_df['记录ID'].astype(str)
if '调仓ID' in processed_df.columns and '股票代码' in processed_df.columns:
processed_df['唯一标识'] = processed_df['股票代码'].astype(str) + '_' + processed_df['调仓ID'].astype(str)
else:
processed_df['唯一标识'] = processed_df['记录ID'].astype(str)
# 从已处理数据中恢复委托状态到原始数据
status_updated_count = 0
for idx, row in raw_df.iterrows():
unique_id = row['唯一标识']
# 在已处理数据中查找相同唯一标识的记录
matching_record = processed_df[processed_df['唯一标识'] == unique_id]
if not matching_record.empty:
# 如果找到匹配记录,使用已处理数据中的委托状态
existing_status = matching_record.iloc[0]['是否委托']
if existing_status == '':
raw_df.loc[idx, '是否委托'] = ''
status_updated_count += 1
self.logger.debug(f'[数据合并] 恢复记录{unique_id}的委托状态为""')
# 删除临时的唯一标识列
raw_df = raw_df.drop('唯一标识', axis=1)
self.logger.info(f'[数据合并] 成功恢复{status_updated_count}条记录的委托状态')
else:
# 如果没有已处理数据,为原始数据添加默认委托状态
if '是否委托' not in raw_df.columns:
raw_df['是否委托'] = ''
self.logger.info('[数据合并] 无已处理数据,为原始数据添加默认委托状态')
df = raw_df # 使用合并后的原始数据
# 4. 如果雪球原始数据不存在,使用已处理的数据文件
elif os.path.exists(self.data_file_path):
df = pd.read_csv(self.data_file_path, encoding='utf-8-sig')
self.logger.info(f'[数据读取] 雪球原始数据不存在,使用已处理数据文件,共{len(df)}条记录')
else:
self.logger.warning('[数据读取] 未找到任何数据文件,创建空DataFrame')
df = pd.DataFrame()
# 5. 数据为空的处理
if len(df) == 0:
self.logger.warning('[数据处理] 数据文件为空,当前无法获取雪球数据,等待下次循环')
return True # 返回True以继续系统运行
except Exception as e:
self.logger.error(f'[数据处理] 读取数据文件失败: {str(e)}')
return False
# 数据完整性验证
is_valid, error_msg = self.validate_data_integrity(df)
if not is_valid:
self.logger.error(f'[数据验证] 数据完整性验证失败: {error_msg}')
return False
self.logger.info('[数据验证] 数据完整性验证通过')
# 2. 数据字段扩展 - 先扩展字段再去重
try:
df = self.extend_data_fields(df)
except Exception as e:
self.logger.error(f'[数据扩展] 数据字段扩展失败: {str(e)}')
return False
# 3. 数据计算处理
try:
df = self.process_data_calculations(df)
except Exception as e:
self.logger.error(f'[数据计算] 数据计算处理失败: {str(e)}')
return False
# 4. 保存更新后的数据 - 包含所有原始数据记录
try:
if not self.save_data_with_backup(df):
self.logger.error('[数据保存] 保存数据文件失败')
return False
self.logger.info(f'[数据保存] 已保存所有{len(df)}条原始数据记录到CSV文件')
except Exception as e:
self.logger.error(f'[数据保存] 保存数据文件异常: {str(e)}')
return False
# 5. 取消数据去重过滤 - 直接处理所有数据
self.logger.info(f'[数据处理] 处理所有{len(df)}条原始数据记录')
# 6. 交易执行 - 只处理当日调仓记录
# 获取当前日期,格式为YYYYMMDD
current_date = datetime.datetime.now().strftime('%Y%m%d')
# 过滤出当日调仓记录且未委托的交易
pending_trades = df[
(df['是否委托'] == '') &
(df['调仓时间'].str.startswith(current_date))
].copy()
if len(pending_trades) == 0:
self.logger.info(f'[交易执行] 没有需要执行的当日调仓交易(当前日期:{current_date}')
return True
self.logger.info(f'[交易执行] 发现{len(pending_trades)}条当日调仓记录需要执行(当前日期:{current_date}')
# 5. 执行交易
successful_trades = 0
failed_trades = 0
for index, row in pending_trades.iterrows():
try:
symbol = row['代码']
follow_shares = row['跟单股数']
trade_direction = row['交易方向']
strategy_name = row['组合名字']
stock_name = row.get('股票名称', '未知')
# 检查跟单股数是否为NAN
if pd.isna(follow_shares):
self.logger.warning(f'[交易执行] {stock_name}({symbol}) 跟单股数为NAN,跳过交易')
continue
volume = int(follow_shares)
if volume <= 0:
self.logger.warning(f'[交易执行] {stock_name}({symbol}) 跟单股数为0,跳过交易')
continue
self.logger.info(f'[交易执行] 准备{trade_direction} {stock_name}({symbol}) {volume}')
# 执行交易
order_id = self.execute_trade_order(symbol, volume, trade_direction, strategy_name)
# 更新交易状态
success = order_id is not None
df = self.update_trade_status(df, index, success)
if success:
successful_trades += 1
else:
failed_trades += 1
except Exception as e:
self.logger.error(f'[交易执行] 处理第{index+1}条交易失败: {str(e)}')
failed_trades += 1
continue
# 保存最终状态
try:
if not self.save_data_with_backup(df):
self.logger.error('[交易执行] 保存交易状态失败')
except Exception as e:
self.logger.error(f'[交易执行] 保存交易状态异常: {str(e)}')
# 交易统计
self.logger.info(f'[交易统计] 成功: {successful_trades}笔, 失败: {failed_trades}')
# 查询并显示成交信息
self.query_and_display_trades()
return True
except Exception as e:
self.logger.error(f'[数据处理] 处理交易数据发生未预期异常: {str(e)}')
return False
def connect_trading_account(self) -> bool:
"""
连接交易账户
Returns:
连接是否成功
"""
try:
if self.qmt is None:
self.logger.error('[账户连接] QMT接口未初始化')
return False
self.qmt._connect()
self.logger.info('[账户连接] 交易账户连接成功')
return True
except Exception as e:
self.logger.error(f'[账户连接] 连接交易账户失败: {str(e)}')
return False
def run_continuous_monitoring(self):
"""
持续监控模式
按照业务逻辑文档要求,在交易时间内持续运行
"""
self.logger.info('[持续监控] 启动持续监控模式')
# 系统健康检查
if not self.validate_system_health():
self.logger.error('[持续监控] 系统健康检查失败,退出程序')
return
retry_count = 0
max_retries = 5
while True:
try:
# 检查是否为交易时间
is_trading, time_msg = self.is_trading_time()
if not is_trading:
self.logger.info('[持续监控] 当前非交易时间,等待中...')
time.sleep(3) # 非交易时间每分钟检查一次
continue
self.logger.info('[持续监控] 开始执行交易业务逻辑')
# 执行主要业务逻辑
success = self.process_trading_data()
if success:
retry_count = 0 # 成功后重置重试计数
self.logger.info('[持续监控] 业务逻辑执行成功')
else:
retry_count += 1
self.logger.warning(f'[持续监控] 业务逻辑执行失败,重试次数: {retry_count}')
if retry_count >= max_retries:
self.logger.error('[持续监控] 达到最大重试次数,暂停监控')
time.sleep(300) # 暂停5分钟
retry_count = 0
# 等待下一次执行
interval = self.config.get('循环间隔', 30)
self.logger.info(f'[持续监控] 等待{interval}秒后继续监控')
time.sleep(interval)
except KeyboardInterrupt:
self.logger.info('[持续监控] 接收到停止信号,退出监控')
break
except Exception as e:
if not self.handle_critical_error('持续监控', e, retry_count):
self.logger.error('[持续监控] 发生不可恢复的错误,退出程序')
break
retry_count += 1
time.sleep(30) # 错误后等待30秒再重试
def run_once(self):
"""
单次运行模式
执行一次完整的交易业务逻辑
"""
self.logger.info('[单次运行] 开始执行单次交易业务逻辑')
try:
# 系统健康检查
if not self.validate_system_health():
self.logger.error('[单次运行] 系统健康检查失败')
return False
# 检查交易时间
is_trading, time_msg = self.is_trading_time()
if not is_trading:
self.logger.warning(f'[单次运行] {time_msg}')
return False
# 执行业务逻辑
success = self.process_trading_data()
if success:
self.logger.info('[单次运行] 交易业务逻辑执行成功')
else:
self.logger.error('[单次运行] 交易业务逻辑执行失败')
return success
except Exception as e:
self.handle_critical_error('单次运行', e)
return False
def main():
"""
主函数入口
根据业务逻辑文档要求实现系统启动逻辑
"""
import sys
try:
# 创建交易系统实例
trading_system = XueqiuTradingSystem()
# 根据命令行参数选择运行模式
if len(sys.argv) > 1 and sys.argv[1] == '--once':
print('[运行模式] 单次运行模式')
success = trading_system.run_once()
sys.exit(0 if success else 1)
else:
print('[运行模式] 持续监控模式(默认)')
print('[运行提示] 可用参数: --once (单次运行)')
trading_system.run_continuous_monitoring()
except KeyboardInterrupt:
print('\n程序被用户中断')
sys.exit(0)
except Exception as e:
print(f'程序启动失败: {str(e)}')
sys.exit(1)
if __name__ == '__main__':
main()