mirror of
https://github.com/silencesdg/mt5_python_ea_suite.git
synced 2026-08-03 22:27:44 +00:00
feat: 一票制并发锁 + MT5硬止损兜底
- 一票制: _pending_long/_pending_short 计数器防同一周期内多信号穿透 - 硬止损: MT5下单时附带sl/tp,倍率1.5x(止损)/1.3x(止盈),比EA软止损更宽 - config.py: 新增 hard_sl_multiplier/hard_tp_multiplier - 全部 send_order 链路(abc/remote/live/server/dryrun/backtest) 支持 sl/tp 参数 - 对冲模块: 彻底移除 - 日志: 去重+30轮摘要
This commit is contained in:
+48
-8
@@ -32,10 +32,12 @@ from config import (
|
||||
SYMBOL, TIMEFRAME, OPTIMIZER_COUNT, OPTIMIZER_START_DATE, OPTIMIZER_END_DATE,
|
||||
USE_DATE_RANGE, INITIAL_CAPITAL, SIGNAL_THRESHOLDS, DEFAULT_WEIGHTS,
|
||||
RISK_CONFIG, GENETIC_OPTIMIZER_CONFIG,
|
||||
MARKET_STATE_CONFIG, TREND_INDICATOR_WEIGHTS, TREND_THRESHOLDS, CONFIDENCE_THRESHOLDS
|
||||
MARKET_STATE_CONFIG, TREND_INDICATOR_WEIGHTS, TREND_THRESHOLDS, CONFIDENCE_THRESHOLDS,
|
||||
DATA_PROVIDER_MODE, REMOTE_SERVER_HOST, REMOTE_SERVER_PORT
|
||||
)
|
||||
from utils.constants import PERIOD_H1
|
||||
from core.utils import get_rates, initialize, shutdown
|
||||
from core.data.remote import RemoteDataProvider
|
||||
from logger import logger
|
||||
|
||||
# 导入策略模块确保 StrategyRegistry 已注册
|
||||
@@ -46,6 +48,18 @@ _multi_tf: MultiTimeframeDataStore | None = None
|
||||
_cached_signals: pd.DataFrame | None = None
|
||||
_registry: StrategyRegistry | None = None
|
||||
|
||||
# MT5 结构化数组 dtype(用于远程 API JSON → numpy 转换)
|
||||
_MT5_RATES_DTYPE = np.dtype([
|
||||
('time', 'i8'),
|
||||
('open', 'f8'),
|
||||
('high', 'f8'),
|
||||
('low', 'f8'),
|
||||
('close', 'f8'),
|
||||
('tick_volume', 'i8'),
|
||||
('spread', 'i4'),
|
||||
('real_volume', 'i8'),
|
||||
])
|
||||
|
||||
|
||||
def init_worker():
|
||||
"""多进程worker初始化:抑制日志噪音"""
|
||||
@@ -397,15 +411,41 @@ def evaluate_fitness(individual, multi_tf: MultiTimeframeDataStore,
|
||||
return (total_pnl,)
|
||||
|
||||
|
||||
def _json_to_mt5_rates(rates_data):
|
||||
"""将远程API JSON rates 转换为 MT5 兼容的 numpy 结构化数组"""
|
||||
if not rates_data:
|
||||
return None
|
||||
records = []
|
||||
for r in rates_data:
|
||||
records.append((
|
||||
r['time'], r['open'], r['high'], r['low'], r['close'],
|
||||
r.get('tick_volume', 0), r.get('spread', 0), r.get('real_volume', 0),
|
||||
))
|
||||
return np.array(records, dtype=_MT5_RATES_DTYPE)
|
||||
|
||||
|
||||
def load_historical_data():
|
||||
"""一次性加载M1数据并返回 MultiTimeframeDataStore"""
|
||||
initialize()
|
||||
rates = (
|
||||
get_rates(SYMBOL, TIMEFRAME, OPTIMIZER_COUNT, OPTIMIZER_START_DATE, OPTIMIZER_END_DATE)
|
||||
if USE_DATE_RANGE else
|
||||
get_rates(SYMBOL, TIMEFRAME, OPTIMIZER_COUNT)
|
||||
)
|
||||
shutdown()
|
||||
if DATA_PROVIDER_MODE == "remote":
|
||||
provider = RemoteDataProvider(host=REMOTE_SERVER_HOST, port=REMOTE_SERVER_PORT)
|
||||
if not provider.initialize():
|
||||
raise RuntimeError("远程MT5 API初始化失败,请检查 Windows MT5 是否运行")
|
||||
|
||||
rates_json = provider.get_historical_data(SYMBOL, TIMEFRAME, OPTIMIZER_COUNT)
|
||||
provider.shutdown()
|
||||
|
||||
if not rates_json:
|
||||
raise RuntimeError("远程获取历史数据失败")
|
||||
|
||||
rates = _json_to_mt5_rates(rates_json)
|
||||
else:
|
||||
initialize()
|
||||
rates = (
|
||||
get_rates(SYMBOL, TIMEFRAME, OPTIMIZER_COUNT, OPTIMIZER_START_DATE, OPTIMIZER_END_DATE)
|
||||
if USE_DATE_RANGE else
|
||||
get_rates(SYMBOL, TIMEFRAME, OPTIMIZER_COUNT)
|
||||
)
|
||||
shutdown()
|
||||
|
||||
if rates is None or len(rates) == 0:
|
||||
raise RuntimeError("获取历史数据失败")
|
||||
|
||||
@@ -3,7 +3,7 @@ import signal
|
||||
import sys
|
||||
from datetime import datetime
|
||||
from logger import logger
|
||||
from config import SYMBOL, TIMEFRAME, REALTIME_CONFIG, SIGNAL_THRESHOLDS
|
||||
from config import SYMBOL, TIMEFRAME, REALTIME_CONFIG, SIGNAL_THRESHOLDS, RISK_CONFIG, RISK_CONFIG_CONST
|
||||
from core.risk import RiskController
|
||||
from execution.weights import DynamicWeightManager
|
||||
|
||||
@@ -16,6 +16,7 @@ class RealtimeTrader:
|
||||
self.running = False
|
||||
self.risk_controller = None
|
||||
self.weight_manager = None
|
||||
self._cycle_count = 0
|
||||
|
||||
def _initialize(self):
|
||||
if not self.data_provider.initialize():
|
||||
@@ -23,25 +24,38 @@ class RealtimeTrader:
|
||||
|
||||
self.risk_controller = RiskController(self.data_provider)
|
||||
self.weight_manager = DynamicWeightManager(self.data_provider)
|
||||
|
||||
self.risk_controller.sync_state()
|
||||
|
||||
signal.signal(signal.SIGINT, self._signal_handler)
|
||||
signal.signal(signal.SIGTERM, self._signal_handler)
|
||||
logger.info("实时交易系统初始化完成")
|
||||
|
||||
# ★ 启动参数一览
|
||||
logger.info("=" * 50)
|
||||
logger.info(f"品种: {SYMBOL} | 周期: M{TIMEFRAME} | 间隔: {self.update_interval}s")
|
||||
logger.info(f"风控: 止损={RISK_CONFIG['stop_loss_pct']:.1%} | "
|
||||
f"止盈={RISK_CONFIG['take_profit_pct']:.1%} | "
|
||||
f"拖尾激活={RISK_CONFIG['min_profit_for_trailing']:.1%} | "
|
||||
f"拖尾回撤={RISK_CONFIG['profit_retracement_pct']:.1%}")
|
||||
logger.info(f"信号: 买入阈值={SIGNAL_THRESHOLDS.get('buy_threshold',1.5)} | "
|
||||
f"卖出阈值={SIGNAL_THRESHOLDS.get('sell_threshold',-1.5)}")
|
||||
logger.info(f"仓位: 最多多={REALTIME_CONFIG['max_long_positions']} 最多空={REALTIME_CONFIG['max_short_positions']} | "
|
||||
f"超时平仓={'开' if RISK_CONFIG_CONST.get('enable_time_based_exit',True) else '关'}")
|
||||
logger.info(f"对冲: 信号对冲={'开' if REALTIME_CONFIG.get('hedge_enabled',False) else '关'} | "
|
||||
f"锁仓={'开' if REALTIME_CONFIG.get('lock_enabled',False) else '关'}")
|
||||
logger.info("=" * 50)
|
||||
return True
|
||||
|
||||
def _signal_handler(self, signum, frame):
|
||||
logger.info(f"接收到信号 {signum},准备退出...")
|
||||
logger.info(f"接收信号 {signum},准备退出...")
|
||||
self.stop()
|
||||
|
||||
def _run_cycle(self):
|
||||
try:
|
||||
self._cycle_count += 1
|
||||
self.risk_controller.sync_state()
|
||||
|
||||
current_price = self.data_provider.get_current_price(SYMBOL)
|
||||
if not current_price:
|
||||
logger.warning("无法获取当前价格,跳过本次循环")
|
||||
return
|
||||
|
||||
strategies_with_weights = self.weight_manager.get_current_strategies_and_weights()
|
||||
@@ -51,66 +65,38 @@ class RealtimeTrader:
|
||||
for strat, weight in strategies_with_weights:
|
||||
signals.append(strat.generate_signal())
|
||||
weights.append(weight)
|
||||
|
||||
# 打印详细信号日志
|
||||
logger.info("--- 信号计算详情 ---")
|
||||
for i, (strat, weight) in enumerate(strategies_with_weights):
|
||||
signal = signals[i]
|
||||
weighted_signal = signal * weight
|
||||
strat_name = strat.name
|
||||
logger.info(f" 策略: {strat_name:<25} | 信号: {signal:6.2f} | 权重: {weight:6.2f} | 加权信号: {weighted_signal:6.2f}")
|
||||
logger.info("--------------------")
|
||||
|
||||
weighted_signal_sum = sum(s * w for s, w in zip(signals, weights))
|
||||
|
||||
buy_threshold = SIGNAL_THRESHOLDS.get('buy_threshold', 1.5)
|
||||
sell_threshold = SIGNAL_THRESHOLDS.get('sell_threshold', -1.5)
|
||||
logger.info(f"加权信号: {weighted_signal_sum:.2f} (买入阈值: {buy_threshold}, 卖出阈值: {sell_threshold})")
|
||||
# logger.info(f"信号比较: {weighted_signal_sum} > {buy_threshold} = {weighted_signal_sum > buy_threshold}")
|
||||
# logger.info(f"信号比较: {weighted_signal_sum} < {sell_threshold} = {weighted_signal_sum < sell_threshold}")
|
||||
|
||||
|
||||
direction = None
|
||||
if weighted_signal_sum > buy_threshold:
|
||||
direction = "buy"
|
||||
elif weighted_signal_sum < sell_threshold:
|
||||
direction = "sell"
|
||||
|
||||
|
||||
# 只在信号触发时打印决策依据
|
||||
if direction:
|
||||
logger.info(f"准备执行{direction}交易,信号强度: {weighted_signal_sum:.2f}")
|
||||
success = self.risk_controller.process_trading_signal(direction, current_price, weighted_signal_sum)
|
||||
if not success:
|
||||
logger.warning(f"{direction}交易执行失败")
|
||||
else:
|
||||
logger.info(f"{direction}交易执行成功")
|
||||
|
||||
logger.info(f"⚡ 信号触发 | 加权={weighted_signal_sum:.2f} | "
|
||||
f"阈值=[{sell_threshold:.2f}, {buy_threshold:.2f}] | "
|
||||
f"方向={direction.upper()} | 价格={current_price['last']:.2f}")
|
||||
self.risk_controller.process_trading_signal(direction, current_price, weighted_signal_sum)
|
||||
|
||||
self.risk_controller.monitor_positions(current_price, weighted_signal=weighted_signal_sum)
|
||||
|
||||
# --- 状态汇总日志 ---
|
||||
logger.info("--- 财务状况更新 ---")
|
||||
open_positions = self.risk_controller.position_manager.positions
|
||||
if not open_positions:
|
||||
logger.info(" 当前无持仓")
|
||||
else:
|
||||
logger.info(f" 当前持仓: {len(open_positions)} 个")
|
||||
for pos in open_positions:
|
||||
pnl_pct = self.risk_controller.position_manager._calculate_pnl_pct(pos, current_price['last'])
|
||||
# 计算持仓时间
|
||||
holding_time = current_price['time'] - pos['entry_time']
|
||||
holding_minutes = holding_time.total_seconds() / 60
|
||||
logger.info(f" - Ticket {pos['ticket']}: {pos['position_type']} {pos['symbol']} @ {pos['entry_price']:.2f} | 持仓时间: {holding_minutes:.1f}分钟 | 浮动盈亏: {pnl_pct:.2%}")
|
||||
|
||||
trade_summary = self.risk_controller.position_manager.get_trade_summary()
|
||||
if trade_summary and trade_summary['total_trades'] > 0:
|
||||
logger.info(" 已平仓交易摘要:")
|
||||
logger.info(f" - 总交易: {trade_summary['total_trades']}, 盈利: {trade_summary['winning_trades']}, 亏损: {trade_summary['losing_trades']}, 胜率: {trade_summary['win_rate']:.2f}%")
|
||||
logger.info(f" - 总净盈亏: ${trade_summary['total_profit_loss']:.2f}")
|
||||
|
||||
logger.info(f" 总权益: ${self.risk_controller.position_manager.total_equity:.2f}")
|
||||
logger.info("----------------------")
|
||||
# 每30个周期打印一次状态摘要
|
||||
if self._cycle_count % 30 == 0:
|
||||
pm = self.risk_controller.position_manager
|
||||
n = len(pm.positions)
|
||||
summary = pm.get_trade_summary()
|
||||
logger.info(f"📊 周期#{self._cycle_count} | 持仓={n} | "
|
||||
f"净值=${pm.total_equity:.2f} | "
|
||||
f"已平{summary['total_trades']}笔 胜率{summary['win_rate']:.0f}% 净${summary['total_profit_loss']:+.2f}")
|
||||
|
||||
except Exception as e:
|
||||
import traceback
|
||||
logger.error(f"交易周期执行失败: {e}\n{traceback.format_exc()}")
|
||||
logger.error(f"交易周期失败: {e}\n{traceback.format_exc()}")
|
||||
|
||||
def start(self):
|
||||
if not self._initialize(): return
|
||||
@@ -132,10 +118,12 @@ class RealtimeTrader:
|
||||
if self.risk_controller:
|
||||
timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
|
||||
self.risk_controller.save_trade_history(f"realtime_trades_{timestamp}")
|
||||
summary = self.risk_controller.position_manager.get_trade_summary()
|
||||
if summary['total_trades'] > 0:
|
||||
logger.info(f"📊 本次运行: {summary['total_trades']}笔 胜率{summary['win_rate']:.0f}% 净${summary['total_profit_loss']:+.2f}")
|
||||
except Exception as e:
|
||||
logger.error(f"保存交易记录失败: {e}")
|
||||
finally:
|
||||
self.data_provider.shutdown()
|
||||
logger.info("实时交易系统已停止")
|
||||
sys.exit(0)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user