QUANTAXIS 数据流处理与事件驱动架构深度解析:从迭代器到分布式消息队列
2026/9/23 21:29:42 网站建设 项目流程
  • 金融科技
  • 后端
  • 数据分析

【免费下载链接】QUANTAXIS

QUANTAXIS 支持任务调度 分布式部署的 股票/期货/期权 数据/回测/模拟/交易/可视化/多账户 纯本地量化解决方案

项目地址:https://gitcode.com/gh_mirrors/qu/QUANTAXIS
点击查看免费下载

导读:本文聚焦 QUANTAXIS 的核心技术底座——数据流(DataFlow)处理与事件驱动架构,完整还原"数据获取 → 逐条推送 → 策略计算 → 下单撮合 → 账户回调 → 绩效分析"的量化闭环。你将掌握panel_gen/secrity_gen迭代器、add_func/add_funcx批量指标叠加、QAEventMQ + QAPubSub + QAThread跨进程分发,以及如何仅切换数据源与账户接口即可让同一策略从回测无缝迁移到实盘。

为什么说"量化的一切都是数据流"

在正式展开回测、实盘等具体场景之前,理解 QUANTAXIS 的数据流思想是提纲挈领的一步。量化交易中遇到的几乎每一个要素,本质上都是一条不断被推送、被消费、再生成新流的数据流

  • 行情类原始流:股票/期货的实时价格数据、实时的买卖盘(Level-2 盘口)数据、实时的财务数据;
  • 派生计算流:基于价格与成交量计算出的指标数据流,基于指标生成的信号流,基于信号生成的订单流;
  • 账户与风控流:基于价格和持仓计算的账户动态权益数据流,基于动态权益计算的绩效与风险数据流;
  • 上层应用流:基于以上所有数据流的可视化,以及基于数据流的预测和分析。

无论你是日内交易者、高频交易者、中长线交易者、纯量化交易、手工量化组合交易、ETF 套利还是 FOF 基金管理,最终比拼的都是对上述数据流的实时处理能力。由此 QUANTAXIS 给出的解决方案核心可以概括为三点:

  1. 数据流的处理能力——迭代器逐条推送、函数批量叠加;
  2. 多数据流的组合能力——指标流/信号流/订单流/账户流互相衔接;
  3. 数据流实时可视化的能力——基于QA_Community界面与QAWebServer的实时展示。

这也是为什么本系列的 P2 阶段要在讲回测之前,先把数据流这个"地基"讲透。

用回测场景直观理解数据流:一条订单的完整旅程

数据流不是抽象概念,QUANTAXIS 的每一次回测都是一次标准的数据流运转。以日线回测为例,整个流程可以拆成清晰的 7 步:

QA_fetch_stock_day_adv --> DataStruct --> panel_gen 迭代器逐条推送(重复循环) | | | ===== 策略部分(on_bar/on_tick)===== | 3.1 缓存历史数据供指标计算 | 3.2 add_func 叠加指标,计算行情信号 | 3.3 依据持仓/现金/权益/风险判断是否下单 | | | ===== 回测引擎部分 ===== | 4. 引擎收到订单 -> 撮合 -> 调用 order.trade 成交回调 | 5. 账户收到回调 -> 更新持仓/现金/权益/风险 | v 当数据流推送完毕 --> 关闭回测 --> QA_Risk / QA_Performance 绩效分析 | v 存储回测产物 --> risk.save / performance.save --> QA_Community 可视化

第 1 步:数据获取,得到 DataStruct

from QUANTAXIS.QAFetch.QAQuery_Advance import QA_fetch_stock_day_adv data = QA_fetch_stock_day_adv('000001', '2020-01-01', '2020-12-31')

QA_fetch_stock_day_adv定义在 QAQuery_Advance.py,返回的是一个多层级索引(date × code)的QA_DataStruct结构体——它是后续一切数据流操作的载体。

第 2 步:DataStruct 迭代器把数据一条条推出来

DataStruct内部把整张面板数据封装成 Python 迭代器,基于yield实现"推送一条、消费一条",数据流由此形成(详见下文迭代器一节)。

第 3 步:策略消费数据流(on_bar / on_tick)

策略在收到on_bar(或on_tick)回调时,面对的是数据流中的当前这一条K 线,策略需要做三件事:

  • 3.1 缓存数据:把历史 K 线缓存下来,方便使用历史数据计算指标;
  • 3.2 搭载指标、计算信号ind = data.add_func(indfunc),把自定义或内置指标函数批量施加到数据流上;
  • 3.3 判断是否下单:基于账户的持仓/现金/权益/风险状态决定是否下单,通过Account.send_order发出订单。在 qactabase.py 中,send_order的签名是send_order(direction='BUY', offset='OPEN', price=3925, volume=10, order_id='', code=None),对应"买卖方向 + 开平仓 + 价格 + 数量"的标准委托要素。

第 4~5 步:回测引擎撮合与账户回调

订单进入回测引擎后:

  • 引擎进行撮合,撮合成功后触发order.trade回调——QAOrder类中trade(trade_id, trade_price, trade_amount, trade_time)记录成交流水(见 QAOrder.py);
  • 账户收到成交回调后调用Account.receive_order,同步更新持仓、现金、权益与风险敞口。

第 6~7 步:收尾分析、存储与可视化

当整个数据流被完全推送结束,回测自动关闭,进入分析阶段:

  • QA_Risk(QA_Account)进行风险分析、QA_Performance(QA_Account)进行收益绩效分析——在 qactabase.py 中即可看到risk = QA_Risk(self.acc)的标准用法;
  • 通过risk.save/performance.save把回测产物持久化存储,随后可在QA_Community界面中做可视化复盘。

关键洞察:这套流程中,策略只关心"数据流里来了什么、我该怎么反应";当策略被放入实盘/模拟环境时,需要改变的仅仅是数据流的来源(从数据库回放换成实时行情推送)以及账户的接口(从模拟撮合换成真实柜台),策略核心逻辑一行都不用改。

QUANTAXIS 数据流处理工具一:迭代器(微型处理单元)

QA_DataStruct提供了两个核心迭代器属性,实现"在不结束数据流的情况下将数据一条一条推送出来":

迭代器按什么维度切片用途
panel_gen时间(date level)切片每次吐出一个"全市场某时刻截面",适合事件驱动回测逐 bar 推进
secrity_gen(文档亦写作security_gen代码(code level)切片每次吐出一个标的的完整序列,适合按标的批量处理

其底层实现位于 base_datastruct.py:

@property def panel_gen(self): '返回一个基于bar的面板迭代器' for item in self.index.levels[0]: # level 0 = 时间 yield self.new( self.data.xs(item, level=0, drop_level=False), dtype=self.type, if_fq=self.if_fq ) @property def security_gen(self): '返回一个基于代码的迭代器' for item in self.index.levels[1]: # level 1 = 代码 yield self.new( self.data.xs(item, level=1, drop_level=False), dtype=self.type, if_fq=self.if_fq )

两者都是惰性生成器(基于yield),因此内存占用与数据总量无关,只与"当前这一条"有关——这正是支撑分钟级、毫秒级长时间回放的关键设计。配套属性bar_gen则以iterrows()形式返回 DataFrame 行迭代器,适合对单行数据做轻量消费。

QUANTAXIS 数据流处理工具二:批量函数叠加 add_func / add_funcx

迭代器解决的是"一条条拿数据",而add_func/add_funcx解决的是"对整条数据流批量施加计算"。这是 QUANTAXIS 数据流处理里最高频的 API:

  • add_func(func, *args, **kwargs):在DataStruct上叠加指标/函数,通过groupby(level=1).apply(func)每个标的分组应用函数,实现"一次调用,多周期多品种批量计算";
  • add_funcx(func, *args, **kwargs):与add_func的区别是会先reset_index变成单索引(pd.DatetimeIndex)再应用函数,适合那些对单层索引 DataFrame 编写的指标函数。

实现见 base_datastruct.py:

def add_func(self, func, *arg, **kwargs): """QADATASTRUCT的指标/函数apply入口""" return self.groupby(level=1, sort=False).apply(func, *arg, **kwargs) def add_funcx(self, func, *arg, **kwargs): """add_funcx 和 add_func 的区别是: add_funcx 会先 reset_index 变成单索引(pd.DatetimeIndex)""" return self.groupby(level=1, sort=False).apply( lambda x: func(x.reset_index(1), *arg, **kwargs))

在指标结构体层面,QAIndicatorStruct.py 也提供了自己的add_func,通过groupby(level=1, as_index=False, group_keys=False).apply(func, raw=True)以 numpy 原始数组模式批量计算,性能更优。典型用法:

# 对多标的日线数据流批量叠加自研指标 data_with_ind = data.add_func(my_indicator_func, *args, **kwargs) # 再叠加第二个指标,形成"指标流 → 信号流"的级联 signal = data_with_ind.add_funcx(signal_func)

正因为add_func面向"多周期、多品种的批量 DataStruct",指标流、信号流可以像流水线一样层层叠加,这正是"多数据流组合能力"的落地形式。

跨进程/分布式的数据流处理:QAEventMQ + QAPubSub + QAThread

单机单进程的迭代器吞吐有限,当数据流需要跨进程、跨机器传递时,QUANTAXIS 提供了完整的分布式数据流方案:

QAEventMQ + quantaxis_pubsub + qathread
  • QAEventMQ:基于消息队列的事件中间件,负责把行情事件、订单事件、账户事件包装成可路由的消息;
  • quantaxis_pubsub(QAPubSub):负责消息的发布/订阅模型,支持广播、路由等多种自由模式,使多个进程/多台机器的多个进程可以互相传递数据和事件;
  • qathread(QAThread/QAEngine 线程体系):负责消费线程的调度与事件循环。

在仓库中对应实现为 QAPubSub(含producer.py/consumer.py/declaters.py等)与 QAEngine(含QAThreadEngine.py/QAAsyncThread.py等)。生产实践中,行情收集进程(如QASU/save_tdx.py)把数据流发布进消息队列,回测/实盘/可视化等消费进程各自订阅所需主题,实现"一份数据流、多端消费"的解耦架构。

这套组合的意义在于:数据流的处理能力从单进程迭代器升级为分布式管道,为后续分布式回测、多账户并行、实时行情分发提供了基础设施。

从回测到实盘:只换数据源与账户接口

数据流架构带来的最大红利是环境可移植性。回测与实盘共享同一套策略代码,差异仅在两处:

环节回测环境模拟/实盘环境
数据流来源QA_fetch_*从本地 Mongo 拉取历史数据回放实时行情推送(TDX/CTP/交易所网关等)
账户接口本地模拟撮合的QA_Account(QIFI 账户结构)真实/模拟柜台接口(QAMarket下的 broker 适配)

从源码结构看,QAStrategy 中的qactabase.pyqamultibase.py等策略基类均面向"数据流 + 账户"抽象编程,而 QAMarket 与 QIFI 分别提供了订单/持仓/账户的统一模型——这种松耦合设计正是"切换环境不切换策略"的保证。

环境搭建:quantaxis_service 一键部署

要快速体验上述全部数据流能力,官方推荐使用quantaxis_service 的 Docker 一键部署:它把 Mongo、消息队列、行情服务等依赖打包成容器,帮助你在任意机器上快速搭建起完整的 QUANTAXIS 运行环境,从而"摆脱环境安装配置的困扰,专注于解决问题本身"。仓库中提供了对应的 compose 编排(见 docker/qa-service 与 docker/qaservice_docker.sh)。

部署完成后,即可按本系列的顺序依次实践:数据准备(P1_Prepare)→ 数据流处理(本文 P2_DataFlow)→ 回测构建(P3_Backtest)→ 绩效分析(P4_Analysis)→ 实时交易(P5_REALTIME)。

小结

数据流是 QUANTAXIS 一切功能的底层逻辑:迭代器解决单机逐条推送,add_func/add_funcx解决批量指标与信号叠加,QAEventMQ + QAPubSub + QAThread解决跨进程分布式分发,而统一的策略/账户抽象让回测与实盘共享同一套数据流消费代码。理解了这条主线,后续无论是构建回测、接入实盘还是搭建可视化,都能做到举一反三。

  • 金融科技
  • 后端
  • 数据分析

【免费下载链接】QUANTAXIS

QUANTAXIS 支持任务调度 分布式部署的 股票/期货/期权 数据/回测/模拟/交易/可视化/多账户 纯本地量化解决方案

项目地址:https://gitcode.com/gh_mirrors/qu/QUANTAXIS
点击查看免费下载

相关推荐

上一篇:ZeroBot-Plugin技术文档自动化:保持文档最新
下一篇:ESLint no-new 规则深度解析:禁止无副作用的 `new` 运算符调用

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询