尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

基于miniQMT构建本地量化数据采集系统:架构设计与工程实践

基于miniQMT构建本地量化数据采集系统:架构设计与工程实践 1. 项目缘起为什么我要自己动手采集miniQMT数据做量化交易的朋友尤其是自己写策略、做回测的应该都懂一个痛点数据。市面上数据源不少有免费的有付费的但总感觉差点意思。免费的像Tushare、AKShare接口稳定性和数据质量有时会抽风而且对调用频率和权限有诸多限制深度数据比如逐笔委托更是难以触及。付费的Wind、iFinD固然好但成本不菲对于个人开发者或者小团队来说是一笔不小的开销。更重要的是很多策略的实盘交易最终要落地到具体的券商交易接口上。如果你的数据源和你的交易接口是两套系统中间就存在一个“对齐”的鸿沟。比如你用A数据源做了回测觉得策略很牛但实盘时通过B券商的接口下单你会发现滑点、成交速度、甚至可交易标的范围都可能和回测环境有差异导致策略失效。这就是所谓的“回测与实盘一致性”问题。我最近在折腾一个“AI股票小助手”的项目核心是想用一些轻量级的AI模型比如时序预测、模式识别辅助交易决策。项目到了第四步数据就成了拦路虎。我需要的是能够紧密贴合我未来实盘交易环境的数据。这时券商的量化交易终端进入了我的视野尤其是miniQMT。miniQMT是某头部券商推出的极速量化交易系统它提供了Python API允许用户直接编写策略并下单延迟极低。但它的数据接口官方文档说得比较含糊社区里也都是零散的代码片段没有一套完整、稳定、可复用的数据采集方案。所以我决定自己动手基于miniQMT的Python API搭建一个本地化的数据采集系统。这个系统的目标很明确稳定、高效、可扩展地获取与实盘交易环境高度一致的行情和基础数据为后续的策略回测和AI模型训练提供高质量的“燃料”。2. 深入理解miniQMT的数据接口能力与限制在动手写代码之前我们必须先摸清楚miniQMT这个“武器库”里到底有什么以及它的“使用说明书”上没写的那些坑。这步没做好后面代码写得再漂亮也是空中楼阁。2.1 官方接口概览我们能拿到什么数据通过研读官方文档和反复测试我梳理出miniQMT Python API主要提供以下几类数据接口实时行情订阅这是核心功能。你可以订阅股票、基金、指数、期货等多个市场的实时tick数据、五档快照、逐笔成交等。数据通过回调函数推送给你的程序延迟可以做到毫秒级对于高频或准高频策略是必备的。历史K线数据可以获取指定标的、周期1分钟、5分钟、日线等、时间范围的历史K线数据。这是回测的基础。但需要注意历史数据的深度和精度不同券商、不同版本的miniQMT可能有差异。基本面与财务数据可以获取股票的基本信息名称、所属板块等、财务指标PE、PB、流通股本等。这部分数据对于基本面量化或因子选股策略很重要。板块与概念数据可以获取行业板块、概念板块的列表以及板块内的成分股。用于做板块轮动或热点追踪策略。交易相关数据如账户资金、持仓、订单状态等。这部分严格来说是交易接口但对于数据采集系统我们主要关注行情和基础数据。2.2 关键限制与“潜规则”那些文档里没明说的事经过实际踩坑我总结了以下几个必须提前知道的限制这直接决定了我们采集系统的架构设计频率限制与流控miniQMT对API调用有严格的频率限制。短时间内高频请求历史K线或基本面数据很容易触发流控导致接口暂时被禁返回空数据或错误。我的经验是对于历史数据拉取必须在请求间加入人工延时如time.sleep(0.5)并且最好以“批次”为单位进行避免连续单次调用。数据完整性校验通过get_market_data或类似函数获取的历史K线有时会存在缺失。比如某只股票在某个交易日停牌返回的日线数据里可能就直接没有这一天的记录而不是用一个NaN或0成交量来占位。这会导致你的时间序列出现断层在后续处理时如计算指标、对齐时间戳带来大麻烦。时间戳的时区与格式miniQMT返回的时间戳可能是本地时间北京时间也可能是UTC时间并且格式可能是整数如20240517150500代表2024年5月17日15:05:00也可能是字符串。必须在代码中统一进行时区转换和格式化强烈建议全部转换为pandas.Timestamp类型并明确时区Asia/Shanghai这是后续所有时间序列操作的基石。合约与标的代码的映射股票代码比较简单如600519.SH但对于基金、可转债、指数代码规则可能不同。而且miniQMT内部可能有一套自己的标的ID体系。你需要建立一个可靠的“外部代码如000001.SZ - miniQMT内部代码”的映射表并处理好退市、更名等情况。连接稳定性miniQMT客户端本身可能因为网络、券商服务器维护等原因断开。我们的采集程序必须具备断线重连机制。不能因为一次网络波动就导致全天数据缺失。注意不同券商对miniQMT的二次开发支持程度不同某些高级数据字段如Level2的买卖队列可能并未完全开放。在开始前最好向你的客户经理或券商的技术支持确认可用接口清单。3. 数据采集系统的核心架构设计基于上述的能力和限制我们不能写一个简单的、一次性的脚本。我们需要一个健壮的、可长期运行的系统。我设计的核心架构如下图所示此处以文字描述整个系统分为四大模块以数据流为主线串联配置与调度中心功能这是系统的大脑。它从一个中心化的配置文件如config.yaml或数据库中读取需要采集的标的列表、数据种类实时tick、历史K线、采集频率、存储路径等参数。设计要点使用schedule或APScheduler这样的库来实现定时任务。例如每天收盘后15:30自动触发“日级历史数据补全”任务交易时间内定时如每5分钟触发“实时快照数据落盘”任务。配置中心要易于修改和扩展。数据获取引擎功能这是与miniQMT API直接交互的模块。它接收调度中心的指令调用相应的miniQMT函数获取数据。核心类设计我会封装一个QMTDataFetcher类。这个类负责初始化并维护与miniQMT客户端的连接实现断线重连逻辑。提供统一的方法如fetch_history_bars(symbol, period, start_time, end_time)来获取历史K线内部处理频率限制和错误重试。提供实时行情订阅的启动/停止方法并设置好数据到达的回调函数。关键技巧在这个引擎内部所有对miniQMT的调用都必须被try...except包裹并记录详细的日志。遇到流控错误通常是特定的错误码应进入指数退避重试逻辑。数据处理与缓存层功能原始数据从API获取后不能直接存储或使用必须经过清洗和格式化。清洗规则去重实时数据可能因网络抖动导致重复推送。异常值处理剔除价格、成交量等字段明显不合理的数据如价格为0或负数。时间戳标准化如前所述统一转换为pandasTimestamp并确保时区正确。字段名标准化将miniQMT返回的字段名可能是中文或缩写映射为你策略中使用的英文标准字段名如‘open‘, ‘high‘, ‘low‘, ‘close‘, ‘volume‘, ‘amount‘。数据对齐对于历史K线检查是否有缺失的交易日并用NaN进行填充保持时间序列的连续性。缓存在数据写入永久存储之前先在内存如Python字典或本地临时文件如Parquet格式中进行缓存。攒够一定数量如1000条或到达一定时间如每10秒再批量写入。这能极大减少对磁盘的I/O操作提升效率。存储与归档模块功能将处理好的数据持久化保存。存储选型实时Tick数据数据量巨大且写入频繁。推荐使用TimescaleDB基于PostgreSQL的时序数据库或DolphinDB。如果追求简单可按日期和标的分片存储为Parquet文件这是一种列式存储格式被pandas和PyArrow完美支持压缩比高读取速度极快。历史K线数据数据量相对较小但查询频繁。可以存储到SQLite或MySQL数据库中方便按条件查询。同样Parquet文件也是极佳选择特别是结合pandas的read_parquet可以快速读取指定时间范围的数据。元数据标的列表、采集任务状态等使用轻量级的SQLite或配置文件即可。目录结构设计一个清晰的目录结构至关重要。例如data/ ├── tick/ # 实时tick数据 │ ├── 2024-05-17/ # 按日期分目录 │ │ ├── 600519.SH.parquet │ │ └── 000001.SZ.parquet │ └── ... ├── kline/ # K线数据 │ ├── 1min/ │ ├── 5min/ │ ├── daily/ │ └── ... ├── fundamental/ # 基本面数据 └── config.yaml # 配置文件4. 核心代码实现与避坑指南接下来我分享几个最核心模块的代码片段和其中踩过的坑。假设我们已经有了一个基本的miniQMT连接对象xt_trader。4.1 历史K线数据获取的稳健实现获取历史数据是最高频的操作也是流控的重灾区。import pandas as pd import time import logging from datetime import datetime, timedelta class QMTDataFetcher: def __init__(self, xt_trader): self.xt_trader xt_trader self.logger logging.getLogger(__name__) # 请求间隔用于规避流控单位秒 self.request_interval 0.5 def fetch_history_bars(self, symbol, period1d, start_time20230101, end_time20231231): 稳健地获取历史K线数据 :param symbol: 标的代码如 600519.SH :param period: 周期1m-1分钟, 5m-5分钟, 1d-日线 :param start_time: 开始时间格式YYYYMMDD或YYYY-MM-DD :param end_time: 结束时间格式同上 :return: pandas.DataFrame, 列包括 [time, open, high, low, close, volume, amount] # 1. 参数转换与验证 # miniQMT可能需要的周期格式转换 period_map {1d: 1d, 1m: 1m, 5m: 5m} qmt_period period_map.get(period) if not qmt_period: raise ValueError(f不支持的周期类型: {period}) # 将时间字符串转换为datetime对象便于后续处理 try: start_dt pd.to_datetime(start_time) end_dt pd.to_datetime(end_time) except Exception as e: self.logger.error(f时间参数格式错误: {e}) return pd.DataFrame() # 2. 分批次获取数据避免单次请求数据量过大或触发流控 all_data [] current_start start_dt # 假设每次最多获取90个交易日的数据根据券商限制调整 delta_batch timedelta(days90) while current_start end_dt: current_end min(current_start delta_batch, end_dt) self.logger.info(f正在获取 {symbol} 从 {current_start.date()} 到 {current_end.date()} 的数据) try: # 调用miniQMT API (此处为示例实际函数名可能不同) # 注意这里的 get_market_data 是示例请替换为实际的API函数名和参数 data self.xt_trader.get_market_data( stock_code[symbol], periodqmt_period, start_timecurrent_start.strftime(%Y%m%d), end_timecurrent_end.strftime(%Y%m%d), field[time, open, high, low, close, volume, amount] ) if data is not None and not data.empty: all_data.append(data) else: self.logger.warning(f获取到空数据: {symbol}, {current_start} to {current_end}) except Exception as e: # 特别处理流控或网络错误 if frequency in str(e).lower() or limit in str(e).lower(): self.logger.warning(f可能触发流控等待10秒后重试: {e}) time.sleep(10) continue # 不增加current_start重试本次批次 else: self.logger.error(f获取数据时发生未知错误: {e}) break # 3. 严格遵守请求间隔 time.sleep(self.request_interval) current_start current_end timedelta(days1) # 下一天开始 # 4. 合并与清洗数据 if not all_data: return pd.DataFrame() df pd.concat(all_data, ignore_indexTrue) df.sort_values(time, inplaceTrue) df.drop_duplicates(subset[time], inplaceTrue) # 基于时间戳去重 # 5. 时间戳标准化 (关键步骤) df[time] pd.to_datetime(df[time]) # 假设原始时间是北京时间但没有时区信息 df[time] df[time].dt.tz_localize(Asia/Shanghai) # 6. 处理可能的缺失日期仅对日线数据 if period 1d: full_date_range pd.date_range(startstart_dt, endend_dt, freqB, tzAsia/Shanghai) # B 工作日 df df.set_index(time).reindex(full_date_range).reset_index() df.rename(columns{index: time}, inplaceTrue) # 对于停牌日成交量、成交额应为0或NaN价格沿用前收这里根据策略决定。 # 简单处理填充NaN # df.fillna(methodffill, inplaceTrue) # 前向填充价格 # df[[volume, amount]] df[[volume, amount]].fillna(0) return df避坑要点分批次这是避免流控和内存溢出的关键。不要试图一次性拉取好几年的数据。错误处理与重试特别是对流控错误必须捕获并采用“指数退避”策略重试。时间戳pd.to_datetime和tz_localize是黄金搭档务必确保整个系统使用统一的时区时间。数据完整性对日线数据做reindex是保证时间序列连续性的好方法但需要谨慎处理停牌日的价格填充逻辑。4.2 实时行情订阅与落盘实时数据追求的是稳定性和低延迟架构设计比代码细节更重要。import threading import queue from pathlib import Path class RealTimeDataCollector: def __init__(self, xt_trader, data_queue, save_interval10): :param xt_trader: miniQMT连接对象 :param data_queue: 线程安全的队列用于缓冲数据 :param save_interval: 批量落盘的时间间隔秒 self.xt_trader xt_trader self.data_queue data_queue self.save_interval save_interval self.is_running False self.subscribed_symbols set() self.logger logging.getLogger(__name__) # 初始化一个字典来缓存不同标的的数据 self.data_buffer {} # symbol - list of data dicts def on_tick_data(self, data): miniQMT行情回调函数 # data 是一个字典包含标的、时间、价格、成交量等信息 symbol data.get(symbol) if not symbol: return # 将数据放入队列由后台线程处理 self.data_queue.put(data) def start_subscription(self, symbol_list): 订阅实时行情 for symbol in symbol_list: if symbol not in self.subscribed_symbols: # 调用miniQMT订阅函数 (示例) # self.xt_trader.subscribe(symbol, callbackself.on_tick_data) self.subscribed_symbols.add(symbol) self.logger.info(f已订阅实时行情: {symbol}) self.is_running True # 启动后台落盘线程 save_thread threading.Thread(targetself._batch_save_worker, daemonTrue) save_thread.start() def _batch_save_worker(self): 后台工作线程定时将缓存中的数据写入文件 while self.is_running: time.sleep(self.save_interval) if not self.data_buffer: continue # 1. 从队列中取出所有累积的数据 data_to_save [] while not self.data_queue.empty(): try: data_to_save.append(self.data_queue.get_nowait()) except queue.Empty: break if not data_to_save: continue # 2. 按标的分组添加到缓存 for data in data_to_save: symbol data.get(symbol) if symbol not in self.data_buffer: self.data_buffer[symbol] [] # 简单的数据清洗确保必要字段存在 cleaned_data { time: pd.to_datetime(data.get(time)).tz_localize(Asia/Shanghai), price: float(data.get(price, 0)), volume: int(data.get(volume, 0)), # ... 其他字段 } self.data_buffer[symbol].append(cleaned_data) # 3. 批量写入Parquet文件 (按日期和标的) self._flush_buffer_to_disk() def _flush_buffer_to_disk(self): 将缓存中的数据写入磁盘 for symbol, data_list in self.data_buffer.items(): if not data_list: continue df pd.DataFrame(data_list) # 按日期创建目录和文件 trade_date pd.Timestamp.now(tzAsia/Shanghai).strftime(%Y-%m-%d) save_dir Path(f./data/tick/{trade_date}) save_dir.mkdir(parentsTrue, exist_okTrue) file_path save_dir / f{symbol}.parquet # 如果文件已存在则读取原有数据并追加 if file_path.exists(): existing_df pd.read_parquet(file_path) df pd.concat([existing_df, df], ignore_indexTrue).drop_duplicates(subset[time]) # 保存 df.to_parquet(file_path, indexFalse) self.logger.debug(f已保存 {len(data_list)} 条tick数据到 {file_path}) # 清空缓存 self.data_buffer.clear()设计精髓生产者-消费者模型回调函数on_tick_data是快速的生产者只负责将数据放入队列。耗时的I/O操作写磁盘由独立的消费者线程_batch_save_worker完成互不阻塞。批量写入定时或定量批量写入磁盘将大量的小IO操作合并为少量的大IO操作性能提升巨大。缓存与去重在内存中按标的缓存数据并定期落盘。落盘前进行去重防止重复数据。4.3 元数据管理标的列表与任务状态一个容易被忽视但至关重要的部分是元数据管理。你需要知道哪些标的的数据已经采集了哪些还没有今天任务执行成功了吗我建议使用一个轻量级的SQLite数据库来管理-- 创建数据库表 CREATE TABLE IF NOT EXISTS collection_tasks ( id INTEGER PRIMARY KEY AUTOINCREMENT, symbol TEXT NOT NULL, data_type TEXT NOT NULL, -- kline_1d, tick, fundamental date DATE NOT NULL, -- 对于历史数据表示数据日期对于任务表示计划采集日期 status TEXT NOT NULL DEFAULT pending, -- pending, success, failed retry_count INTEGER DEFAULT 0, created_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP, updated_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP, UNIQUE(symbol, data_type, date) -- 防止重复任务 );你的调度中心在每天启动时会查询这个表找出所有statuspending且date为今天或之前的任务然后交给数据获取引擎去执行。执行成功后更新状态为success失败则增加retry_count超过阈值后标记为failed并发出告警。这套机制能保证采集任务的幂等性和可追溯性。5. 系统部署、监控与日常运维代码写完了让它7x24小时稳定运行才是真正的挑战。5.1 部署方式选择本地Windows/Mac电脑最简单但受限于电脑开关机和不稳定。云服务器推荐购买一台Linux云服务器如腾讯云、阿里云。你需要解决在无图形界面的Linux服务器上运行miniQMT客户端的问题。这通常需要用到xvfb虚拟显示来模拟一个图形环境。# 示例使用xvfb-run来启动需要图形界面的程序 xvfb-run --auto-servernum --server-num1 --server-args-screen 0 1024x768x24 python your_data_collector.pyDocker容器化进阶将整个采集程序、依赖库、甚至包括通过xvfb运行的miniQMT客户端打包成一个Docker镜像。这可以实现一键部署和环境隔离是更专业的做法。5.2 监控与告警系统不能黑盒运行。你需要知道它是否还活着数据是否在正常采集。心跳监控在采集程序中每隔一段时间如5分钟向一个监控文件或数据库写入一条“心跳”记录。用一个独立的监控脚本定期检查这个心跳。如果超过一定时间如10分钟没有更新就认为程序僵死。日志监控使用logging模块将日志同时输出到控制台和文件。重点监控ERROR和WARNING级别的日志。可以使用ELKElasticsearch, Logstash, Kibana栈或更轻量的LokiGrafana来集中管理和告警。数据质量监控每天收盘后运行一个检查脚本。例如检查今天应该采集的标的数量 vs 实际采集到的标的数量。检查每个标的的数据是否有缺失例如日线数据是否缺少了某一天。检查数据的极端值如价格跳变超过10%。告警渠道当监控到异常时通过邮件、企业微信、钉钉机器人或Telegram Bot发送告警信息给你让你能及时介入处理。5.3 日常运维清单每日开盘前登录服务器检查miniQMT客户端是否正常登录检查日志文件是否有前一日未处理的错误。每日收盘后触发数据完整性校验脚本确保当日数据完整无误。备份当日数据到冷存储如OSS、AWS S3。每周/每月检查磁盘空间清理过期的临时文件。对数据库进行优化如SQLite的VACUUM命令。回顾监控告警记录分析系统薄弱环节并进行优化。从头搭建一个基于miniQMT的数据采集系统就像在给自家的量化交易引擎建造一个专属的、高质量的“油库”。这个过程充满了挑战从理解晦涩的API文档到处理各种边界情况和异常流控再到设计高可用的系统架构。但一旦这个系统稳定运行起来你获得的将是一份与你的实盘交易环境同源、干净、结构化的高质量数据集。这份数据才是你后续进行策略回测、AI模型训练、乃至实盘交易的真正底气。它让你摆脱了对第三方数据源的依赖也让你的策略研究闭环更加完整和可靠。
返回列表