在实际量化交易开发中,我们常常面临一个核心矛盾:策略研究需要 Python 的灵活生态和强大库支持,而交易执行则依赖于券商或平台提供的专用终端(如迅投 QMT、恒生 PTrade)。这些终端通常内置了策略编写环境,但功能受限,难以使用外部的 Python 库进行复杂的数据处理和模型计算。一个常见的需求是,能否让在外部独立运行的、功能强大的 Python 策略程序,与 QMT 这样的交易终端进行实时、双向的通信,从而实现“外部研究,内部执行”的架构?
这就是信号桥接方案要解决的问题。本文将深入探讨一种基于文件通信的稳健桥接方案。该方案不依赖网络端口,不涉及复杂的进程间通信(IPC)库,而是利用操作系统文件系统作为消息队列,实现外部 Python 程序与 QMT 终端之间的指令与数据交换。这种方案部署简单,跨平台兼容性好,且易于调试和监控,非常适合作为量化交易系统中的一个可靠通信层。
本文将带你从零开始,理解桥接的核心思想,设计通信协议,并实现一个包含外部 Python 信号发生器与 QMT 内部信号监听器的完整可运行案例。你将掌握如何让 Python 策略生成交易信号并写入文件,以及如何让 QMT 定时读取文件、解析信号并执行订单。我们还会详细讨论文件锁、信号去重、异常处理等生产环境必须考虑的关键细节。
1. 理解基于文件通信的桥接架构与核心挑战
在深入代码之前,必须厘清基于文件通信的核心工作模式及其背后的工程考量。这种方案的本质是将文件系统的一个特定目录(或文件)充当为一个简单的、持久化的消息队列。
1.1 核心工作流程:生产者-消费者模型
整个系统遵循典型的生产者-消费者模型:
- 生产者(外部 Python 程序): 持续运行的计算引擎。它从数据源(如数据库、网络 API)获取数据,运行复杂的策略模型(可能用到
pandas,numpy,scikit-learn,TA-Lib等 QMT 不直接支持的库),生成交易信号(如:“买入 600519.SH 100股”、“平仓所有 IF 空头头寸”)。然后,它将信号按照约定格式写入一个特定的“信号文件”。 - 消费者(QMT 终端内的脚本): 在 QMT 中运行一个定时任务(例如每秒或每 500 毫秒)。这个任务负责检查“信号文件”是否有新内容。一旦发现新信号,便立即读取、解析,并调用 QMT 的订单接口执行交易。执行完成后,通常需要更新文件状态(如删除已处理信号或写入确认标记),防止信号被重复执行。
1.2 为什么选择文件通信?优势与妥协
选择文件而非网络套接字、消息中间件(如 Redis、RabbitMQ)或数据库,是基于特定场景的权衡:
优势:
- 零依赖: 无需在 QMT 端安装任何额外的 Python 包或中间件。QMT 的 Python 环境通常是封闭且受限的,文件操作是它的基础功能。
- 简单可靠: 文件系统是操作系统最稳定的抽象之一。写入文件是一个原子性相对较高的操作(尤其是在配合文件锁时),比管理网络连接更简单。
- 易于调试: 信号以明文形式保存在文件中,开发人员可以直接查看文件内容来验证信号是否正确生成,极大简化了调试流程。
- 跨平台: 在 Windows、Linux 或 macOS 上,文件操作的 API 基本一致,方案可移植性强。
妥协与挑战:
- 延迟: 文件 I/O 的速度远低于内存或网络通信,不适合超高频(微秒级)交易。但对于秒级、分钟级的策略完全足够。
- 并发控制: 必须妥善处理“边读边写”的冲突,否则会导致信号丢失或重复执行。这是本方案需要解决的核心技术问题。
- 状态管理: 需要设计一套简单的协议来标识信号的状态(如“待处理”、“已处理”、“执行失败”)。
1.3 关键设计决策:通信协议与文件格式
通信协议定义了双方对话的语言。一个健壮的协议至少包含以下要素:
- 信号格式: 推荐使用JSON。它结构清晰,易于人类阅读和机器解析,且 QMT 的 Python 环境原生支持
json库。{ "signal_id": "order_20240520_093001_001", "timestamp": "2024-05-20 09:30:01.500", "action": "BUY", "symbol": "600519.SH", "price": 1650.50, "volume": 100, "strategy_name": "dual_thrust" } - 文件组织方式:
- 单文件滚动: 始终向同一个文件(如
signal.json)追加新信号。消费者读取后,可以清空文件或移动文件。风险在于清空时可能丢失正在写入的信号。 - 文件队列: 每个信号生成一个独立文件(如
signal_20240520093001500.json),消费者处理完后删除或移走该文件。这种方式更安全,但需要管理文件清理。 - 状态文件: 使用两个文件。一个
signal_pending.json用于写入新信号,消费者读取后,将信号移至signal_processed.json或直接删除。 本文将采用“单文件追加 + 文件锁 + 状态标记”的组合方案,在简单性和安全性之间取得平衡。
- 单文件滚动: 始终向同一个文件(如
2. 环境准备与项目结构
在开始编码前,需要明确双方的环境和项目目录结构。
2.1 外部 Python 环境配置
你的外部策略程序可以在任何独立的 Python 环境中运行。确保安装你策略所需的所有库。
# 创建一个新的虚拟环境(推荐) python -m venv venv_qmt_bridge # 激活虚拟环境 # Windows: venv_qmt_bridge\Scripts\activate # Linux/Mac: source venv_qmt_bridge/bin/activate # 安装常用量化分析库(示例) pip install pandas numpy ta-lib scikit-learn2.2 QMT 端环境认知
QMT 终端内置了 Python 解释器。你需要了解:
- Python 版本: 通常是 Python 3.6 或 3.7。你的脚本语法需兼容。
- 工作目录: QMT 脚本的运行目录。通常是你策略脚本所在的目录,或者 QMT 的安装目录下的某个位置。这是桥接方案的关键,双方必须能访问同一个目录或通过网络共享访问同一个目录。
- 可用库: QMT 内置了
json,time,os,sys等标准库,以及它自己的交易 API(如xtquant)。不要假设可以pip install。
2.3 共享目录设置
这是桥接方案的物理基础。两个程序必须能读写同一个目录下的文件。
- 本地开发: 最简单的方式是指定本地磁盘的一个绝对路径,如
D:\qmt_signal_bridge。确保 QMT 有权限读写该目录。 - 生产部署: 如果外部 Python 程序运行在另一台服务器上,则需要使用网络共享文件夹(SMB/NFS)。将共享文件夹映射到 QMT 所在机器的驱动器(如
Z:\),双方都使用这个网络路径。
项目结构示例:
qmt_signal_bridge/ ├── shared_dir/ # 共享目录,双方均可读写 │ ├── signal.json # 信号文件(核心通信媒介) │ └── signal.lock # 文件锁(可选,用于并发控制) ├── external_strategy/ # 外部策略程序 │ ├── strategy.py # 主策略逻辑 │ ├── signal_writer.py # 信号写入模块 │ └── requirements.txt └── qmt_script/ # QMT 内部脚本 ├── bridge_consumer.py # 信号监听与执行主脚本 └── utils.py # 工具函数3. 实现外部 Python 信号生产者
外部策略程序负责生成信号并写入共享文件。这里的关键是安全写入,避免产生损坏或部分写入的信号。
3.1 信号生成与封装
首先,定义一个信号生成函数。在实际项目中,这里会包含复杂的数据获取和策略计算。
# external_strategy/strategy.py import json import time from datetime import datetime from typing import Dict, Any import random def generate_trading_signal() -> Dict[str, Any]: """ 模拟生成一个交易信号。 实际应用中,这里会接入行情数据,运行策略模型。 """ # 模拟一些策略逻辑 symbols = ["600519.SH", "000858.SZ", "300750.SZ"] actions = ["BUY", "SELL", "CANCEL"] signal = { "signal_id": f"signal_{int(time.time() * 1000)}_{random.randint(1000, 9999)}", # 唯一ID "timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S.%f")[:-3], "action": random.choice(actions), "symbol": random.choice(symbols), "price": round(random.uniform(100, 200), 2), "volume": random.randint(100, 10000), "strategy_name": "example_moving_average", "status": "PENDING" # 初始状态:待处理 } return signal if __name__ == "__main__": # 测试信号生成 for _ in range(3): sig = generate_trading_signal() print(json.dumps(sig, indent=2, ensure_ascii=False)) time.sleep(0.5)3.2 安全写入信号到共享文件
直接向文件追加 JSON 对象会遇到问题:文件内容可能不是合法的 JSON 数组。我们采用“读取-更新-写入”模式,并引入文件锁(fcntl在 Windows 上不可用,这里使用portalocker第三方库,或采用原子重命名策略)。
方案一:使用 portalocker 进行跨平台文件锁(推荐用于外部程序)
# 在外部策略环境中安装 portalocker pip install portalocker# external_strategy/signal_writer.py import json import os import time import portalocker from typing import List, Dict, Any class SignalFileWriter: def __init__(self, file_path: str): self.file_path = file_path # 确保文件存在且初始化为空列表 if not os.path.exists(self.file_path): with open(self.file_path, 'w', encoding='utf-8') as f: json.dump([], f) def append_signal(self, signal: Dict[str, Any]) -> bool: """ 线程/进程安全地向信号文件追加一个信号。 返回是否成功。 """ max_retries = 3 for attempt in range(max_retries): try: # 使用文件锁打开文件 with open(self.file_path, 'r+', encoding='utf-8') as f: portalocker.lock(f, portalocker.LOCK_EX) # 获取排他锁 # 读取现有信号列表 try: existing_signals = json.load(f) except json.JSONDecodeError: # 如果文件损坏,则重置 existing_signals = [] # 追加新信号 existing_signals.append(signal) # 写回文件 f.seek(0) # 移动到文件开头 f.truncate() # 清空文件内容 json.dump(existing_signals, f, indent=2, ensure_ascii=False) # 锁会在 with 块结束时自动释放 print(f"[Writer] Successfully appended signal: {signal['signal_id']}") return True except (IOError, OSError, json.JSONDecodeError) as e: print(f"[Writer] Attempt {attempt + 1} failed: {e}") time.sleep(0.1) # 短暂等待后重试 print(f"[Writer] Failed to write signal after {max_retries} attempts.") return False def read_all_signals(self) -> List[Dict[str, Any]]: """读取当前文件中所有信号(通常用于调试)。""" try: with open(self.file_path, 'r', encoding='utf-8') as f: portalocker.lock(f, portalocker.LOCK_SH) # 获取共享锁 return json.load(f) except (FileNotFoundError, json.JSONDecodeError): return []方案二:原子写入策略(不依赖第三方库)如果不想引入portalocker,可以采用“写入临时文件,然后原子替换”的方式,这在大多数文件系统上是安全的。
# external_strategy/signal_writer_atomic.py import json import os import tempfile from typing import Dict, Any, List class SignalFileWriterAtomic: def __init__(self, file_path: str): self.file_path = file_path self.temp_dir = os.path.dirname(file_path) def append_signal(self, signal: Dict[str, Any]) -> bool: max_retries = 3 for attempt in range(max_retries): try: # 1. 读取现有内容 existing_signals = [] if os.path.exists(self.file_path): try: with open(self.file_path, 'r', encoding='utf-8') as f: existing_signals = json.load(f) except json.JSONDecodeError: existing_signals = [] # 文件损坏则重置 # 2. 追加新信号 existing_signals.append(signal) # 3. 写入临时文件 with tempfile.NamedTemporaryFile( mode='w', dir=self.temp_dir, delete=False, encoding='utf-8', suffix='.tmp' ) as tmp_file: tmp_path = tmp_file.name json.dump(existing_signals, tmp_file, indent=2, ensure_ascii=False) # 4. 原子替换原文件 (Windows 下可能需要重试) # os.replace() 在支持上是原子的 os.replace(tmp_path, self.file_path) print(f"[Writer] Successfully appended signal atomically: {signal['signal_id']}") return True except (IOError, OSError) as e: print(f"[Writer] Attempt {attempt + 1} failed: {e}") # 清理可能残留的临时文件 if 'tmp_path' in locals() and os.path.exists(tmp_path): try: os.unlink(tmp_path) except OSError: pass time.sleep(0.1) return False3.3 整合策略循环与信号写入
将信号生成器和写入器组合起来,形成一个持续运行的生产者。
# external_strategy/main_producer.py import time import signal as sys_signal import sys from strategy import generate_trading_signal from signal_writer import SignalFileWriter # 或 SignalFileWriterAtomic class ProducerDaemon: def __init__(self, signal_file_path: str, interval_seconds: float = 5.0): self.writer = SignalFileWriter(signal_file_path) self.interval = interval_seconds self.running = True # 优雅退出处理 sys_signal.signal(sys_signal.SIGINT, self.signal_handler) sys_signal.signal(sys_signal.SIGTERM, self.signal_handler) def signal_handler(self, signum, frame): print(f"\n[Producer] Received signal {signum}, shutting down...") self.running = False def run(self): print(f"[Producer] Started. Writing signals to {self.writer.file_path} every {self.interval}s") while self.running: try: # 1. 生成信号 new_signal = generate_trading_signal() # 2. 写入文件 success = self.writer.append_signal(new_signal) if not success: print("[Producer] Warning: Failed to write last signal.") except Exception as e: print(f"[Producer] Error in main loop: {e}") # 等待下一个周期 for _ in range(int(self.interval * 10)): # 以0.1秒为粒度检查,便于响应退出信号 if not self.running: break time.sleep(0.1) if __name__ == "__main__": # !!!重要:这里需要修改为实际的共享文件路径 !!! SHARED_SIGNAL_FILE = r"D:\qmt_signal_bridge\shared_dir\signal.json" # 或者网络路径,如 r"Z:\shared_signal\signal.json" daemon = ProducerDaemon(SHARED_SIGNAL_FILE, interval_seconds=5) daemon.run() print("[Producer] Stopped.")运行此生产者,它将在指定路径的signal.json文件中持续追加交易信号。
4. 实现 QMT 端信号消费者
QMT 端的脚本需要定时读取信号文件,解析并执行交易,然后更新信号状态。
4.1 QMT 脚本基础框架与定时任务
在 QMT 的 Python 策略编辑器中,通常使用xtquant等模块。我们首先搭建一个定时扫描的框架。
# qmt_script/bridge_consumer.py # -*- coding: utf-8 -*- """ QMT信号桥接消费者主脚本。 将此脚本在QMT中创建为策略并运行。 需要配置:共享信号文件的正确路径。 """ import json import os import time import traceback from datetime import datetime # QMT相关模块 import xtquant.xtdata as xt_data import xtquant.xttrader as xt_trader import xtquant.xttype as xt_type class QmtSignalConsumer: def __init__(self, signal_file_path: str, scan_interval_ms: int = 1000): """ 初始化消费者。 :param signal_file_path: 共享信号文件的绝对路径。 :param scan_interval_ms: 扫描文件的时间间隔(毫秒)。 """ self.signal_file_path = signal_file_path self.scan_interval = scan_interval_ms / 1000.0 # 转换为秒 self.processed_ids = set() # 用于去重,防止重启后重复执行 self._ensure_file_exists() # 初始化QMT交易接口(示例,根据实际API调整) # self.session_id = None # self._init_trader() def _ensure_file_exists(self): """确保信号文件存在,如果不存在则创建空文件。""" if not os.path.exists(self.signal_file_path): try: with open(self.signal_file_path, 'w', encoding='utf-8') as f: json.dump([], f) print(f"[QMT Consumer] Created new signal file at {self.signal_file_path}") except IOError as e: print(f"[QMT Consumer] ERROR: Cannot create signal file: {e}") # def _init_trader(self): # """初始化QMT交易接口,需要根据券商和实际API填写。""" # try: # # 示例代码,具体API请参考迅投QMT官方文档 # # self.session_id = xt_trader.create_session() # # xt_trader.login(session_id, ...) # print("[QMT Consumer] Trader API initialized (placeholder).") # except Exception as e: # print(f"[QMT Consumer] ERROR initializing trader API: {e}") def _read_signals_safely(self): """ 安全地读取信号文件。 使用文件锁或原子读取策略,避免读取到不完整的JSON。 注意:QMT环境可能没有portalocker,这里使用简单的重试机制。 """ signals = [] for _ in range(3): # 重试3次 try: if not os.path.exists(self.signal_file_path): return [] with open(self.signal_file_path, 'r', encoding='utf-8') as f: content = f.read().strip() if not content: return [] signals = json.loads(content) # 基本验证:确保读到的是列表 if isinstance(signals, list): return signals else: print(f"[QMT Consumer] WARNING: File content is not a list: {type(signals)}") return [] except json.JSONDecodeError as e: print(f"[QMT Consumer] JSON decode error (可能文件正在写入),重试... 错误: {e}") time.sleep(0.05) # 等待50毫秒后重试 except IOError as e: print(f"[QMT Consumer] IOError reading file: {e}") break return [] # 重试后仍失败,返回空列表 def _process_single_signal(self, signal: dict): """ 处理单个信号。 1. 检查是否已处理(去重)。 2. 解析信号内容。 3. 调用QMT交易接口下单。 4. 标记信号为已处理。 """ signal_id = signal.get('signal_id') if not signal_id: print(f"[QMT Consumer] ERROR: Signal missing 'signal_id': {signal}") return False # 去重检查 if signal_id in self.processed_ids: print(f"[QMT Consumer] Signal {signal_id} already processed, skipping.") return True # 视为成功,但跳过执行 action = signal.get('action', '').upper() symbol = signal.get('symbol', '') volume = signal.get('volume', 0) price = signal.get('price', 0.0) print(f"[QMT Consumer] Processing signal: id={signal_id}, action={action}, symbol={symbol}, volume={volume}") # 根据信号执行交易(这里是核心交易逻辑) success = self._execute_trade(action, symbol, volume, price) if success: # 记录已处理ID,防止重启后重复执行(注意:内存记录,重启会丢失) self.processed_ids.add(signal_id) # 更可靠的做法:将已处理信号ID持久化到另一个文件 self._mark_signal_processed(signal_id) print(f"[QMT Consumer] Successfully processed signal {signal_id}") else: print(f"[QMT Consumer] Failed to execute signal {signal_id}") return success def _execute_trade(self, action: str, symbol: str, volume: int, price: float) -> bool: """ 调用QMT交易接口执行订单。 此处为示例伪代码,必须替换为真实的QMT API调用。 参考 xtquant.xttrader 模块的文档。 """ # !!! 重要:以下代码需要根据迅投QMT实际API修改 !!! try: # 示例伪代码: # if action == "BUY": # order = xt_type.StockOrder(...) # order.symbol = symbol # order.volume = volume # order.price = price # order.order_type = xt_type.LIMIT_ORDER # order.side = xt_type.BUY # result = xt_trader.order(self.session_id, order) # elif action == "SELL": # ... 类似 ... # else: # print(f"Unsupported action: {action}") # return False # return result and result.success print(f"[QMT Consumer] (Placeholder) Would execute: {action} {volume} of {symbol} @ {price}") # 模拟执行成功 time.sleep(0.1) # 模拟网络延迟 return True except Exception as e: print(f"[QMT Consumer] ERROR in trade execution: {e}") traceback.print_exc() return False def _mark_signal_processed(self, signal_id: str): """ 将信号标记为已处理。 简单实现:从当前信号列表中移除已处理的信号,并写回文件。 更复杂的实现可以使用独立的处理日志文件。 """ try: current_signals = self._read_signals_safely() # 过滤掉已处理的信号(根据ID) pending_signals = [s for s in current_signals if s.get('signal_id') != signal_id] # 写回文件(这里存在并发风险,但在QMT单线程循环中问题不大) # 生产环境建议使用更安全的方式,如写入另一个“已处理”文件。 with open(self.signal_file_path, 'w', encoding='utf-8') as f: json.dump(pending_signals, f, indent=2, ensure_ascii=False) print(f"[QMT Consumer] Removed processed signal {signal_id} from file.") except Exception as e: print(f"[QMT Consumer] ERROR marking signal processed: {e}") def run_forever(self): """主循环,定时扫描并处理信号。""" print(f"[QMT Consumer] Started. Monitoring {self.signal_file_path} every {self.scan_interval}s") while True: # QMT策略通常在一个无限循环中运行 try: # 1. 读取信号 signals = self._read_signals_safely() if signals: print(f"[QMT Consumer] Found {len(signals)} signal(s) to process.") # 2. 处理每个信号 for sig in signals: self._process_single_signal(sig) except Exception as e: print(f"[QMT Consumer] ERROR in main loop: {e}") traceback.print_exc() # 3. 等待下一次扫描 time.sleep(self.scan_interval) # ========== 在QMT中运行的部分 ========== def main(): # !!! 关键配置:必须与生产者写入的路径一致 !!! # 如果是Windows本地,使用如 r'D:\qmt_signal_bridge\shared_dir\signal.json' # 如果是网络映射驱动器,使用如 r'Z:\shared_signal\signal.json' SHARED_SIGNAL_FILE = r"D:\qmt_signal_bridge\shared_dir\signal.json" consumer = QmtSignalConsumer(SHARED_SIGNAL_FILE, scan_interval_ms=2000) # 每2秒扫描一次 consumer.run_forever() if __name__ == "__main__": # 在QMT中,这个脚本通常被作为一个策略加载,main()函数会被自动或手动调用。 # 也可以将QmtSignalConsumer实例化并嵌入到QMT的策略模板中。 main()4.2 在 QMT 中部署与运行脚本
- 创建策略: 在 QMT 的策略编辑器中,新建一个 Python 策略。
- 粘贴代码: 将上述
bridge_consumer.py的核心逻辑(主要是QmtSignalConsumer类和main函数)复制到策略编辑器中。 - 修改配置:
- 将
SHARED_SIGNAL_FILE变量修改为与外部生产者完全一致的共享文件路径。 - 根据你的券商和 QMT 版本,修改
_execute_trade方法内的下单代码。这部分代码必须参考迅投 QMT 的官方开发文档,使用正确的 API 和参数。
- 将
- 设置定时运行: 在 QMT 策略设置中,将策略的运行周期设置为“定时”,间隔可以设为 1 秒或更长(根据你的信号频率调整)。这会让 QMT 定时调用策略的
main函数或相应的处理函数。 - 运行策略: 启动策略。观察 QMT 的输出窗口(日志),应该能看到类似
“[QMT Consumer] Started...”和发现信号、处理信号的日志。
5. 运行验证与调试流程
桥接方案涉及两个独立进程,调试需要系统化。
5.1 验证步骤
- 启动生产者: 在外部 Python 环境中运行
main_producer.py。检查共享目录下是否生成了signal.json文件,并且内容在不断增加新的 JSON 对象。 - 启动消费者: 在 QMT 中运行部署好的策略。观察 QMT 日志,确认它成功读取到了
signal.json文件。 - 观察交互: 在生产者运行一段时间后,查看 QMT 日志是否打印出
“Processing signal...”和“Successfully processed...”的信息。同时,观察signal.json文件,已被处理的信号应该被移除(根据我们的_mark_signal_processed逻辑)。 - 模拟交易: 在 QMT 的模拟交易环境中运行此策略,验证
_execute_trade中的下单逻辑是否能成功报单(即使不真实成交)。
5.2 关键调试点与日志
在两个程序中加入详细的日志是排查问题的关键。
| 环节 | 可能的问题 | 调试方法 |
|---|---|---|
| 文件路径 | 双方访问的不是同一个文件。 | 在双方代码中打印文件的绝对路径os.path.abspath(file_path)进行比对。检查网络映射驱动器的连接状态。 |
| 文件权限 | QMT 或外部程序没有写入/读取权限。 | 检查文件属性,确保运行 QMT 的用户和运行外部 Python 程序的用户都有权限。在共享文件夹设置中检查共享权限。 |
| JSON 格式 | 文件内容不是合法的 JSON,导致解析失败。 | 直接打开signal.json文件查看内容。生产者应确保每次写入都是完整的 JSON 数组。使用本文提供的带锁或原子写入的SignalFileWriter。 |
| 信号重复执行 | QMT 重启后,内存中的processed_ids清空,导致重复执行旧信号。 | 实现持久化的已处理信号记录,例如将已处理的signal_id写入另一个文件processed_ids.log,启动时加载。或者在信号中增加时间戳,消费者忽略过时的信号。 |
| QMT API 错误 | _execute_trade中的下单 API 调用失败。 | 仔细阅读 QMT 官方文档,使用正确的函数、参数和订单类型。先在 QMT 中编写一个最简单的下单脚本测试 API 是否可用。查看 QMT 的错误输出。 |
| 网络延迟 | 网络共享文件夹延迟高,导致信号传递慢。 | 检查网络状况。考虑将扫描间隔 (scan_interval_ms) 适当调大,避免过于频繁的 I/O。 |
6. 生产环境进阶考量与最佳实践
上述方案是一个可运行的最小原型。投入生产前,必须加强其健壮性和可维护性。
6.1 增强通信可靠性
- 双文件乒乓缓冲:
- 使用两个文件:
signal_active.json和signal_backup.json。 - 生产者总是写入
signal_active.json。 - 消费者读取
signal_active.json后,立即将其重命名为signal_backup.json并处理。同时,生产者下次写入时会发现active文件消失,则创建新的active文件。这种方式彻底避免了读写冲突。
- 使用两个文件:
- 信号确认与重试:
- 消费者处理完信号后,不仅从待处理列表移除,还应向一个
ack_<signal_id>.json文件写入确认信息。 - 生产者可以定期扫描
ack_*文件,如果某个信号长时间未被确认,可以重新发送(需注意幂等性)。
- 消费者处理完信号后,不仅从待处理列表移除,还应向一个
6.2 完善信号协议
- 信号状态管理: 在信号对象中增加
status字段,如PENDING,PROCESSING,SUCCESS,FAILED。消费者读取后先改为PROCESSING,执行成功或失败后更新状态。这需要更复杂的文件更新逻辑。 - 信号过期: 增加
expiry_time字段。消费者如果发现信号已过期,则直接丢弃并记录日志,避免执行无效指令。 - 策略上下文: 信号中可以携带策略运行时的快照信息,如
data_hash,model_version,便于后续复盘和审计。
6.3 QMT 脚本的健壮性
- 异常隔离: 确保处理一个信号的异常不会导致整个监控循环崩溃。
_process_single_signal内部要有完善的try...except。 - 资源清理: 如果使用了 QMT 的交易会话,在策略停止时要有清理逻辑。
- 心跳与监控: QMT 脚本可以定期向一个
heartbeat.txt文件写入当前时间。外部监控程序可以检查这个文件,如果超过一定时间未更新,则报警。 - 配置外置: 将文件路径、扫描间隔等配置项放在独立的
config.json中,而不是硬编码在脚本里。
6.4 部署与运维清单
在将本方案部署到实盘前,请逐项检查:
- [ ]共享目录: 确认生产服务器与 QMT 所在机器之间的网络共享稳定,且有冗余链路。
- [ ]文件路径: 双方程序使用的都是绝对路径,且指向同一个物理文件。
- [ ]权限: 生产环境运行程序的用户(可能是系统服务账户)拥有共享目录的读写权限。
- [ ]错误处理: 生产者的写入失败、消费者的读取失败和下单失败都有日志记录和报警通知(如邮件、钉钉、Telegram)。
- [ ]日志轮转: 制定了日志文件(包括信号文件本身如果作为日志)的清理策略,避免磁盘写满。
- [ ]信号去重: 实现了基于持久化存储(如小数据库或文件)的信号去重,防止 QMT 重启后重复下单。
- [ ]熔断机制: 当连续出现多次通信失败或下单失败时,程序应能自动暂停并报警,而不是持续产生错误。
- [ ]回滚方案: 准备了手动清理信号文件、停止双方程序的应急预案。
基于文件通信的 QMT 信号桥接方案,以其简单、可靠、无额外依赖的特性,在分钟级及以上频率的量化策略中是一个非常实用的集成模式。它成功地将策略研究的灵活性与交易执行的稳定性解耦。实现的重点不在于复杂的通信框架,而在于对文件读写并发控制、信号协议设计和异常处理细节的把握。从本文的最小可行方案出发,结合你的具体交易逻辑和风险控制要求,逐步完善其周边设施,可以构建出一个支撑实盘交易的稳健桥梁。