Level 2 量化因子开发和模拟交易的 DolphinDB 最佳实践
在高频金融交易领域,逐笔成交数据和 Level 2 快照行情包含丰富的时序特征与市场微观结构信息,是量化因子构建和交易策略设计的核心基础。随着市场结构日益复杂、算法交易占比持续提升,实时处理海量高频数据的能力已成为量化团队的核心竞争力。高频交易系统每日需处理亿级数据,流程涵盖预处理、时间对齐、分钟频聚合、因子计算及模型推理等环节,对实时性、稳定性和扩展性要求极高。
针对上述需求,DolphinDB 推出了分钟频量化交易实时因子计算模块 SnapTradeStream。该模块基于高性能流计算引擎,实现从实时行情采集、分钟频聚合、因子计算、模型推理、信号生成到模拟交易撮合的全链路处理,具备低延迟和高吞吐能力,有效提升了实时因子的计算效率与系统稳定性。
1. 基于快照行情生成分钟频因子
Level 2 快照行情包含十档买卖盘口、委托量和实时成交价格统计等信息,能反映市场深度、盘口压力及供需失衡信号。然而,原始快照数据每秒可达数千条且字段复杂,直接用于策略逻辑计算负担大、难以提取有效的分钟频特征。本章介绍如何将高频快照数据实时聚合为分钟频因子,包括数据特征、计算规则及实现步骤。
1.1 股票 Level 2 快照行情数据特征
快照数据是交易所每隔固定时间推送的全市场静态行情信息,广泛应用于分钟频行情归档、因子计算和量化策略回测研究。它包含丰富的字段结构,不仅包括基础行情(如最新价、最高价、最低价、成交量等),还涵盖委托队列(多档买卖价和对应委托量)、ETF 申赎数据、买卖盘统计、成交笔数等信息。快照数据字段详见附录。基于这些字段,可计算出以下常用的分钟频因子指标:
| 字段名称 | 数据类型 | 数据说明 |
|---|---|---|
| securityId | SYMBOL | 证券代码 |
| rseComputeUs | LONG | 响应式状态引擎计算时间(单位为微秒) |
| batchSize | INT | 响应式状态引擎计算每批条数 |
| snapTseComputeUs | LONG | 时序聚合引擎计算时间(单位为微秒) |
| tradeTime | TIMESTAMP | 当前聚合窗口时间戳(下一分钟起始) |
| market | SYMBOL | 交易市场 |
| high | DOUBLE | 最高成交价 |
| open | DOUBLE | 开盘价 |
| low | DOUBLE | 最低成交价 |
| close | DOUBLE | 收盘价 |
| volume | LONG | 当前分钟新增成交量 |
| amount | DOUBLE | 当前分钟新增成交额 |
| numTrades | LONG | 当前分钟新增成交笔数 |
| totalSellVolume | LONG | 当前十档总委卖量 |
| totalBuyVolume | LONG | 当前十档总委买量 |
| sellVolume | LONG | 十档卖盘总委托量 |
| buyVolume | LONG | 十档买盘总委托量 |
| entrustSellAmount | DOUBLE | 当前委卖总金额 |
| entrustBuyAmount | DOUBLE | 当前委买总金额 |
| sellAvgPrice | DOUBLE | 当前时点加权平均卖价 |
| buyAvgPrice | DOUBLE | 当前时点加权平均买价 |
1.2 分钟频因子指标计算规则
将高频快照数据转化为分钟频指标,既能降低数据维度,又保留市场核心结构信息,是后续模型推理和信号触发的重要输入。DolphinDB 通过时序聚合引擎与响应式状态引擎实现高效低延迟的分钟频指标计算。以下从价格、成交、盘口委托三个维度介绍计算规则。
1.2.1 价格计算规则
分钟频级别的价格走势是量化策略识别趋势、判断强弱、识别拐点的基本输入。例如,收盘价决定模型的预测目标构建,开盘价与最高/最低价则影响回撤、突破策略的判断。同时,需要合理处理异常(如开盘价为 0)与缺失数据。因此,为保证策略使用的数据一致性与稳定性,分钟频价格指标遵循如下计算逻辑:
-
开盘价(open):
-
对非开盘、收盘分钟,直接取窗口内首条快照最新价字段的值作为开盘价。
-
对开盘和收盘分钟:
-
如果窗口内最新的开盘价为 0,则使用本窗口首条快照最新成交价;
-
否则采用窗口内最新的开盘价。
-
-
-
最高价(high)、最低价(low):
-
如果窗口内最后一条数据的最高价(最低价)为 0,分别取本窗口内最新成交价的最大(最小)值。
-
否则为窗口内最新的最高价(最低价)。
-
-
收盘价(close):
-
取当前窗口最后一条数据的最新价即为收盘价。
-
1.2.2 成交量、成交额和成交笔数
成交指标能够反映市场活跃度,是判断“有效突破”与“虚假拉升”的核心参考。例如,量能放大通常伴随趋势延续,量能萎缩可能预示回调。特别是在盘中择时策略中,成交量的实时变化往往决定是否执行信号。由于快照中成交量是累加字段,为得到准确的“本分钟新增成交”,采用差分方式提取每分钟真实变动:
-
取窗口内最新的一条数据的值暂时作为快照成交量、成交额、成交笔数。
-
使用
deltas()函数分别计算相邻分钟差值,若差值为空则保留窗口内最新的一条数据的值,否则使用计算出的差值。 -
字段包括:
-
volume:新增成交量。
-
amount:新增成交额。
-
numTrades:新增成交笔数。
-
1.2.3 委托量及均价计算规则
快照数据的一大重要价值,在于提供十档盘口的挂单信息,反映当前市场的供需结构。量化策略尤其关注委托量变化,例如通过盘口量价不对称,判断是否存在“压单吸筹”或“拉高出货”。因此,分钟频委托统计是捕捉主力行为的重要窗口。为确保策略捕捉的是每分钟结束时的最新盘口状态,按以下指标对挂单信息进行处理:
-
在计算委托买卖量和金额相关指标时,所有指标均取当前计算窗口内最新一条快照行情的数据:
-
totalSellVolume:委卖总量,取窗口中 totalOfferQuant 字段的最新值。
-
totalBuyVolume:委买总量,取窗口中 totalBuyQuant 字段的最新值。
-
entrustSellAmount:委卖总金额,由卖盘挂单量乘以加权平均卖价计算得到,取窗口内最新一条数据。
-
entrustBuyAmount:委买总金额,类似 entrustSellAmount,由买盘挂单量乘以加权平均买价计算得到。
-
sellAvgPrice,buyAvgPrice:加权平均卖价和买价。
-
1.3 分钟频因子指标计算步骤
分钟频因子指标的构建不仅依赖于精确的计算规则,更依赖于流数据的高质量处理与窗口计算机制。在实盘环境中,原始快照数据存在交易时段复杂、时间戳不规范、数据分布不均等问题,这对分钟频计算的准确性与实时性提出挑战。为了保障分钟频因子稳定、及时地输出,DolphinDB 通过快照预处理 、时序聚合 、响应式状态引擎构建了高性能的分钟频数据链路。下面将分三个阶段介绍 DolphinDB 如何从原始快照数据高效生成分钟频因子数据。
1.3.1 快照数据清洗与时间标准化处理
快照数据中包含集合竞价、连续竞价、不规则时间戳与跨市场数据,若不清洗将严重干扰分钟频因子的准确性。特别是在分钟线触发临界点,不标准的时间戳会导致聚合误差与分钟频数据缺失。因此,需对原始快照进行时间清洗和格式化处理,确保其落入正确的聚合窗口。我们采用如下预处理逻辑对快照数据进行筛选与时间标准化:
-
筛选交易时段内的数据:保留集合竞价时段(09:25–09:30)和连续竞价时段(09:30–11:30,13:00–14:57)内的快照数据,剔除无效或非交易时段的数据。
-
优化分钟线触发时机:考虑到每分钟快照数据的最新交易时间为 59 秒,针对这一特点进行了优化处理——当接收到 59 秒的快照数据时,会立即触发分钟线的计算。
以下代码基于 DolphinDB3.00.4版本开发,建议用户使用3.00.4及以上版本
// tse: 时序聚合引擎
// msg: 原始快照数据
def handleRawSnapshot(mutable tse, mutable msg){
msg = select * from msg
where ((second(tradeTime) between 09:25:00:11:30:59) || (second(tradeTime) between 12:59:00:14:56:59))
if (size(msg) == 0) return
update msg set tradeTime =
case
when minute(tradeTime) in 09:25m..09:29m then concatDateTime(date(tradeTime), 09:30:00.000)
else tradeTime
end
msg1 = select * from msg where second(tradeTime) % 60 >= 59
update msg set tradeTime = timestamp(concatDateTime(date(tradeTime), second(minute(tradeTime)))) + 58 * 1000 + 999 where second(tradeTime) % 60 >= 59
update msg1 set tradeTime = timestamp(concatDateTime(date(tradeTime), second(minute(tradeTime)))) + 59 * 1000 + 1
append!(msg, msg1)
msg.sortBy!(`tradeTime)
append!(tse, msg)
}
1.3.2 时序聚合引擎生成分钟频数据
快照数据经过清洗后,需在每分钟的时间窗口内进行高效聚合,从而输出标准化的分钟频行情。在量化场景下,不同股票的快照数量差异极大,传统聚合方式容易漏算、错算、重复计算。而 DolphinDB 的时序聚合引擎滑窗聚合,确保每分钟窗口在数据到齐后立即触发计算,具体处理逻辑包括:
-
设置窗口与步长:窗口大小为 60 秒(1 分钟),步长也是 60 秒,确保每条分钟线独立且无重叠。
-
强制触发机制:针对分钟窗口内只有少量快照数据的情况,设置forceTriggerTime参数强制输出分钟线。
tse = createTimeSeriesEngine(
name= tseName,
windowSize=60000,
step=60000,
metrics=tseMetrics,
dummyTable=rawStream,
outputTable=rse,
timeColumn=`tradeTime,
useSystemTime=false,
keyColumn=`securityID,
useWindowStartTime=false,
parallelism=parallelTSE,
lowLatency = ctx["lowLatency"],
outputElapsedMicroseconds = true
)
1.3.3 响应式引擎生成分钟频因子
分钟频数据虽然为标准行情摘要,但无法直接用于策略决策,必须进一步加工出如主力行为、盘口强弱、成交偏离等“策略可用特征”。这些因子的生成过程依赖状态变量、滑动窗口、历史状态记忆,需要在毫秒级延迟下完成复杂计算逻辑。为此,DolphinDB 提供响应式状态引擎以流式计算分钟频因子。上一步的分钟聚合结果将被输入至响应式状态引擎,实时生成策略所需的分钟频因子数据:
rse = createReactiveStateEngine(
name = rseName,
metrics = rseMetrics,
dummyTable = tseOutputDummy,
outputTable = objByName(minStreamTbname, true),
keyColumn = `securityId,
parallelism = parallelRSE,
lowLatency = ctx["lowLatency"],
outputElapsedMicroseconds = true
)
2. 基于逐笔成交生成分钟频因子
逐笔成交数据记录了市场上每一笔交易的真实细节,是洞察资金流向、识别大单行为和市场微观结构的重要依据。量化策略尤其依赖对主力资金行为、资金流方向和盘口异动的实时判断。但逐笔数据量巨大(单日可达数亿条),且具有时间不连续性的特点。本章将聚焦于如何利用 DolphinDB 强大的流处理能力,对原始逐笔成交数据进行实时清洗、排序、窗口聚合,精准生成反映分钟频资金动向的关键因子,为捕捉市场情绪和主力行为提供数据支撑。
2.1 逐笔成交数据特征
与快照数据不同,逐笔数据具有更高的时间分辨率,能完整记录每笔成交的价格、数量、方向、成交时间、成交方式等信息。逐笔成交数据字段详见附录。基于这些字段,本文计算出以下常用的分钟频因子指标:
| 字段名 | 数据类型 | 含义 |
|---|---|---|
| securityId | SYMBOL | 证券代码 |
| rseComputeUs | LONG | 响应式状态引擎计算时间(单位为微秒) |
| batchSize | INT | 响应式状态引擎计算每批条数 |
| tradeTseComputeUs | LONG | 时序聚合引擎计算时间(单位为微秒) |
| tradeTime | TIMESTAMP | 当前聚合窗口的时间戳 |
| buyAmount | DOUBLE | 所有买入成交金额之和 |
| sellAmount | DOUBLE | 所有卖出成交金额之和 |
| buySmallAmount | DOUBLE | 过去 1 分钟内,买方向小单的成交额,成交股数小于等于 50,000 股 |
| buyBigAmount | DOUBLE | 过去 1 分钟内,买方向大单的成交额,成交股数大于 50,000 股 |
| sellSmallAmount | DOUBLE | 过去 1 分钟内,卖方向小单的成交额,成交股数小于等于 50,000 股 |
| sellBigAmount | DOUBLE | 过去 1 分钟内,卖方向大单的成交额,成交股数大于 50,000 股 |
| TradeBias | DOUBLE | 最后一笔成交趋势 |
2.2 分钟频因子指标计算规则
DolphinDB 基于时序窗口与响应式引擎,对逐笔成交数据构建了以下两类分钟频因子指标。
2.2.1 成交金额类指标
量化策略通常以资金流向作为判断市场强弱、主力方向的重要输入。而逐笔数据提供了最精细的资金流明细,可以精确拆解为主动买入/卖出金额,进而判断主力行为是否在买入护盘或卖出打压。此外,通过分析最后一笔成交相对于大/小单的偏移,还可以构造微观层面的行为偏差因子。基于这一业务需求,以下指标用于衡量当前分钟内成交的资金流动方向和规模,计算逻辑如下:
-
sellAmount:主动卖出成交金额。
-
buyAmount:主动买入成交金额。
-
TradeBias:衡量最后成交是否偏向大单买入或卖出的因子,计算逻辑为:(lastAmount - sellMinAmount) / abs(sellMinAmount) - (lastAmount - buyMinAmount) / abs(buyMinAmount)
-
lastAmount:当前分钟最后一笔成交金额。
-
buyMinAmount:主动买入方向中,所有单笔成交金额的最小值。
-
sellMinAmount:主动卖出方向中,所有单笔成交金额的最小值。
-
2.2.2 主动资金大/小单拆分指标
在市场博弈中,大资金往往通过大单执行交易,而小单则多为散户行为。将主动成交划分为大单/小单,有助于识别是否有主力资金正在悄悄建仓或出货。特别是在某些个股活跃时段,大单持续出现常常意味着背后存在机构行为,值得策略重点关注。因此,为了刻画不同规模主动资金行为,定义以下大/小单资金流指标,用于分钟频上量化资金行为的结构:
将单笔成交量是否大于 50,000 作为判断标准:
-
buySmallAmount:主动买入且成交量 ≤ 50,000 的资金总额。
-
buyBigAmount:主动买入且成交量 > 50,000 的资金总额。
-
sellSmallAmount:主动卖出且成交量 ≤ 50,000 的资金总额。
-
sellBigAmount:主动卖出且成交量 > 50,000 的资金总额。
2.3 分钟频因子指标计算步骤
逐笔成交数据是交易撮合的最小单位,记录了市场每一笔真实的成交行为。相比快照数据,逐笔数据更能体现真实资金流动与交易行为结构。但其特点是数据量更大、时间精度更细、乱序情况更频繁,因此,在进入分钟频聚合流程前,必须进行专门设计的数据处理与窗口聚合策略。以下将介绍 DolphinDB 如何通过数据标准化 → 时序聚合引擎 → 响应式状态引擎,构建逐笔数据驱动的分钟频实时因子计算流程。
2.3.1 成交数据的预处理与标准化
在逐笔成交数据中,存在乱序、时间格式不统一、跨市场混合等情况。若直接进行聚合计算,容易导致数据遗漏或窗口划分错误,直接影响分钟频因子的精度与完整性。因此,为确保逐笔成交数据在分钟频聚合时准确划分窗口,确保所有成交数据有序,需对原始成交数据进行基本的清洗和时间标准化处理。具体步骤如下:
-
排序对齐:首先对逐笔成交数据按 tradeTime 升序排序,确保数据的时间顺序一致,避免因乱序导致被时序引擎丢弃,造成分钟频指标正确性问题。
-
数据写入:将时间频率转成秒频,写入时序引擎,以供后续分钟频聚合计算使用。
// tse: 时序聚合引擎
// msg: 原始快照数据
def handleRawTrade(mutable tse, mutable msg){
msg.sortBy!(`tradeTime)
// 转成秒频数据,避免乱序问题
replaceColumn!(msg, `tradeTime, timestamp(datetime(msg.tradeTime)))
tse.append!(msg)
}
2.3.2 时序聚合引擎生成分钟频数据
经过标准化的成交数据会被注入时序聚合引擎,进行分钟频窗口划分与初步聚合。逐笔成交密度高但分布不均,必须设置合理的窗口配置与强制触发机制,确保每个自然分钟都有输出。DolphinDB 的时序聚合引擎通过设定parallelism参数,支持多股票并行滑窗聚合,并提供交易时段控制与性能输出功能,确保实盘环境下的实时性。具体处理逻辑如下:
-
窗口配置:设置windowSize=60000和step=60000,即每个聚合窗口大小为 1 分钟,滑动步长也是 1 分钟,实现不重叠的自然分钟划分。
-
强制触发控制:设置forceTriggerTime=100,及时触发不活跃股票的分钟频计算。
tse = createTimeSeriesEngine(
name= tseName,
windowSize=60000,
step=60000,
metrics=tseMetrics,
dummyTable=rawStream,
outputTable=rse,
timeColumn=`tradeTime,
useSystemTime=false,
keyColumn=`securityID,
useWindowStartTime=false,
parallelism=parallelTSE,
lowLatency = ctx["lowLatency"],
outputElapsedMicroseconds = true
)
2.3.3 响应式引擎生成分钟频因子
基础聚合数据仅提供成交总额、笔数等粗粒度信息,无法满足策略对资金结构、行为偏差、大单交易等特征的需求。DolphinDB 通过响应式状态引擎,对分钟频成交数据流进行细粒度加工,构建具有策略解释力的因子特征,并保持实时输出能力。分钟频聚合结果被输入至状态引擎,实现分钟频因子的流式计算与实时输出:
rse = createReactiveStateEngine(
name = rseName,
metrics = rseMetrics,
dummyTable = tseOutputDummy,
outputTable = objByName(minStreamTbname, true),
keyColumn = `securityId,
parallelism = parallelRSE,
lowLatency = ctx["lowLatency"],
outputElapsedMicroseconds = true
)
3. 基于分钟频因子生成衍生因子
基础分钟频因子反映了市场瞬时状态,但单一指标难以捕捉复杂动态和潜在趋势。量化策略通常需要技术指标、统计特征或多因子组合等衍生特征,直接在原始数据上计算效率低且逻辑复杂。此外,衍生因子的计算需将多个异构数据流实时对齐合并,对引擎设计和容错机制要求较高。本章介绍如何基于已生成的分钟频因子表,通过 DolphinDB 流计算引擎进行多因子融合与计算,实时生成策略所需的衍生因子。
3.1 衍生因子定义
本文基于快照分钟频因子与逐笔成交分钟频因子,实时计算以下衍生因子指标:
| 字段名 | 数据类型 | 含义 |
|---|---|---|
| securityId | SYMBOL | 证券代码 |
| tradeTime | TIMESTAMP | 交易分钟时间 |
| dealToEntrustRatio | DOUBLE | 成交额与委托比率因子 |
| bigOrderVsDepth | DOUBLE | 大单成交与盘口深度因子 |
| priceDeviation | DOUBLE | 成交价格偏离度因子 |
| priceImpactEfficiency | DOUBLE | 价格冲击效率因子 |
| bigOrderSupport | DOUBLE | 大单价格支撑因子 |
| tradeMatchEfficiency | DOUBLE | 盘口成交匹配度因子 |
| priceDiscoveryEfficiency | DOUBLE | 价格发现效率因子 |
| liquidityConsumption | DOUBLE | 流动性消耗因子 |
| capitalEfficiency | DOUBLE | 资金效率因子 |
| marketPressure | DOUBLE | 综合市场压力因子 |
3.2 衍生因子指标计算步骤
快照分钟频因子和成交分钟频因子各自刻画了市场的不同侧面,但在实际策略设计中,单一来源的信息往往不足以支撑精细化的交易决策。例如,只有快照指标很难刻画主力买入行为,只有成交指标也难反映盘口压力。因此,需要将不同来源的分钟频特征进行融合,并在此基础上构建更高层次的“衍生因子”。为实现这一过程,DolphinDB 通过 EquiJoinEngine 引擎对多源数据进行分钟频匹配,并使用响应式状态引擎实时计算复杂的复合特征。以下是完整的衍生因子计算链路:
3.2.1 引擎结构与数据流向
衍生因子需要同时利用快照因子和成交因子进行计算,这要求系统具备跨表合并、并发处理和状态维护能力。同时,生成的衍生因子还要供下游推理模块实时调用,必须保证数据结构统一、处理路径清晰。因此,DolphinDB 设计了如下几类核心组件构建数据流处理链路:
-
EquiJoin 流引擎(EquiJoinEngine):用于将两个来源(快照分钟频数据与成交分钟频数据)在 tradeTime 和 securityId 维度上匹配,形成分钟频的合并输入流。
-
响应式状态引擎(ReactiveStateEngine):接收 EquiJoinEngine 输出的合并流数据,在此基础上完成衍生因子计算。该引擎按证券代码进行状态维护与指标更新,支持并发处理。
-
衍生因子流数据表:用于记录状态引擎输出的分钟频因子流,支持持久化、订阅和后续推理生成交易信号。
3.2.2 EquiJoin 引擎合并分钟频数据流
为了使分钟频因子具备更强的表达能力,需将来源不同的分钟频数据进行统一对齐。DolphinDB 的 EquiJoinEngine 提供了毫秒级别的时间对齐能力,确保合并后的数据可直接用于下游状态引擎的衍生因子计算。该引擎对接来自快照与成交的分钟频因子表,对来自两个来源的分钟频数据表进行合并。
-
合并数据源:
-
左表(LeftTable):基于快照行情生成的分钟频因子表。
-
右表(RightTable):基于逐笔成交生成的分钟频因子表。
-
-
合并方式:
-
匹配字段为 securityId。
-
对齐的时间字段为 tradeTime。
-
ejEngine = createEquiJoinEngine(
name = "EquiJoinEngine",
leftTable = snapMinDummy,
rightTable = tradeMinDummy,
outputTable = rse,
metrics = sqlCol(colNames[2:]),
matchingColumn = `securityId,
timeColumn = `tradeTime
)
3.2.3 响应式状态引擎计算分钟频衍生因子
合并后的分钟频数据拥有更丰富的信息密度,可以构造如“成交额偏离盘口均价”、“主力行为 VS 挂单反应”、“大单主导趋势”等更具策略价值的复合因子。这类因子通常涉及复杂计算逻辑,因此使用 DolphinDB 的响应式状态引擎实现高效处理。
rse = createReactiveStateEngine(
name = rseName,
metrics = facMetrics,
dummyTable = dummy,
outputTable = objByName(factorTbname, true),
keyColumn = `securityId,
parallelism = parallelRSE,
lowLatency = ctx["lowLatency"],
outputElapsedMicroseconds = true
)
计算结果最终写入衍生因子流数据表中,供下一步模型实时推理并生成交易信号。
4. 基于衍生因子训练模型实时推理
量化交易的核心在于预测。生成衍生因子后,如何在毫秒级延迟内完成实时预测并输出买卖信号,是实盘策略成败的关键。传统批处理或高延迟推理流程会错失交易机会。本章介绍在 DolphinDB 流处理框架内,如何基于实时衍生因子完成多模型加载与推理,生成交易信号供后续执行。
4.1 训练模型准备
在量化交易系统中,模型的训练和准备是实现精准预测的基础环节。模型训练需从历史数据中提取有效特征,结合市场微观结构构建机器学习模型。训练环节涵盖数据清洗、特征工程和因子选择等步骤,并通过多模型集成综合不同维度信息,实现多视角信号预测。本文准备了 PyTorch 和 LightGBM 两个 Demo 模型用于后续实时推理。
4.2 实时推理步骤
实时推理需在数据产生后毫秒级内完成信号生成,以捕捉短暂交易机会。本文利用 DolphinDB 流计算框架,将因子计算与多模型推理紧密集成,支持 PyTorch 和 LightGBM 模型加载至内存进行实时推理。主要步骤如下:
-
数据接收与预处理:接收流式计算生成的衍生因子数据,进行清洗与格式转换,确保符合各模型输入规范。
-
多模型调用与推理:将预处理数据分别输入多个模型,各模型推理结果独立输出,整合后推送至信号流数据表供下游订阅。
5. 基于 DolphinDB 的实时仿真交易
实时推理生成信号后,需通过仿真交易验证策略的有效性和执行稳健性,在实际资金投入前发现潜在风险。本章介绍如何基于 DolphinDB 搭建实时仿真交易系统,复现市场行情与订单撮合环节,实现从信号生成到委托执行的闭环验证。
注意:实时仿真交易依赖MatchingEngineSimulator和OrderManagementEngine插件,需提前安装与加载。
5.1 仿真交易搭建步骤
仿真交易系统模拟真实市场环境下的订单撮合与成交过程,在风险可控前提下评估策略表现与执行效率。具体步骤如下:
-
交易参数设置与模拟撮合引擎创建:配置交易品种、时间周期、撮合规则及风控参数,初始化模拟撮合引擎。该引擎支持订单簿管理及价格优先、时间优先撮合原则。
-
实时订阅快照行情:订阅原始快照行情流表,获取最新价格与买卖挂单深度,经规范化处理后输入撮合引擎。
-
基于衍生因子生成信号并提交委托订单:订阅衍生因子推理输出的交易信号,结合策略逻辑生成委托订单,实时推送至撮合引擎进行动态匹配。
-
实时撮合与成交回报:撮合引擎依据价格优先、时间优先规则实时撮合买卖委托并生成成交明细,结果持续输出至流数据表,形成信号到成交的闭环。
5.2 模拟撮合引擎与订单管理
在基于 DolphinDB 的实时仿真交易系统中,除了利用模拟撮合引擎重现市场订单撮合机制外,还引入了订单管理的功能。下面将针对模拟撮合引擎订单管理的作用与优势进行详细说明。
5.2.1 时间窗口控制的交易执行
订单管理功能允许用户为每笔订单指定开始时间(effectiveTime)和结束时间(expireTime),只有在这一区间内,订单才能被正常撮合执行,一旦达到结束时间,未成交或部分成交的订单会触发撤单操作。
-
保障交易执行的时效性:防止委托单无限期“挂单”,避免在市场行情已经变化且策略意图改变后订单仍在市场中暴露风险。
-
合理模拟实际交易环境:真实交易中,很多订单都会设置有效期(限时有效),通过订单管理功能,仿真交易更贴合真实约束。
-
提升仿真精度和风控能力:及时清理过期订单,有效控制交易风险,避免因过期订单引发的误判。
5.2.2 部分成交及到期处理机制
在仿真环境中,部分成交的订单可能会剩余未成交数量,订单管理功能负责:
-
按结束时间自动触发剩余订单状态更新,避免残留委托干扰后续策略演化。
-
输出完整的成交和订单状态信息,保证仿真数据的完整性和准确性。
5.2.3 适用场景
-
量化交易策略中,订单生命周期极短,订单管理插件精准模拟订单有效期至关重要。
-
资金管理与风险控制严苛的机构策略,需对订单执行的时限和成交情况设定严格监控。
-
做市商策略中的订单主动撤单和时效控制,保证撮合过程的流畅性和实时反馈。
-
策略测试与回测阶段,对订单撮合时效与部分成交行为进行详细追踪和分析。
6. 实时因子生成及后续推理与模拟撮合模块
构建实时因子计算与预测流水线涉及多个引擎、流表订阅、数据合并和模型调用,手动配置繁琐且易出错,难以满足生产环境的稳健性和扩展性要求。本章介绍基于 DolphinDB 模块机制封装的SnapTradeStream模块,实现从分钟频因子计算、衍生因子生成到模型推理与模拟撮合的完整流程,支持快速部署、灵活配置和历史数据回放,降低了构建实时量化分析平台的运维成本。
下图为实时流计算全流程的整体架构:
6.1 流计算框架结构及功能
SnapTradeStream模块的文件结构为:
SnapTradeStream/
├── CreateTable.dos
├── OHLCFactors.dos
├── OrcaFrame.dos
├── StreamFrame.dos
├── StreamUtils.dos
└── utils.dos
其中,每个文件的功能分别为:
| 文件名 | 功能 |
|---|---|
| utils.dos | 定义模块中使用的通用工具函数 |
| StreamUtils.dos | 定义传统流计算和 Orca 共同用到的工具函数 |
| CreateTable.dos | 定义创建相关数据库表函数 |
| OHLCFactors.dos | 定义分钟频行情因子的字段构造逻辑与因子表达式函数 |
| StreamFrame.dos | 定义从原始数据到衍生指标,再到推理与模拟仿真交易的计算流程函数 |
| OrcaFrame.dos | 定义 Orca 计算平台中从原始数据到衍生指标,再到推理与模拟仿真交易的计算流程函数 |
注:
-
使用 Orca 功能时,DolphinDB 版本需升级至 3.00.3 及以上
-
若使用集群模式,需部署计算组,Orca 调用脚本在计算组其中一个计算节点执行
6.2 流计算框架构建步骤
实时计算系统的核心在于全链路框架的稳定性与低延迟能力,尤其在市场剧烈波动时更为关键。本节介绍 SnapTradeStream 模块构建实时因子计算框架的核心流程,涵盖从原始数据到分钟频聚合、多源合并、衍生因子计算直至推理与模拟撮合的全链路。流程概括如下:
-
分钟频因子处理:将逐笔成交与快照数据分别送入时序聚合引擎(60 秒滚动窗口,基于 tradeTime)和响应式状态引擎,输出标准化分钟频因子流表,并对数据进行结构对齐与字段补充。
-
衍生因子计算:通过 EquiJoinEngine 以 securityId 和 tradeTime 为连接字段合并快照与成交分钟频数据,结果送入响应式状态引擎,生成衍生因子并写入流数据表。
-
推理与模拟仿真:订阅衍生因子流表,传入推理模型生成交易信号,据此构建委托订单并传入模拟撮合引擎,实时输出成交详情。
6.3 流计算框架应用实践
本节将介绍如何通过SnapTradeStream模块从零搭建以上章节介绍的流计算框架,详细代码参考附件文件。
6.3.1 导入成交逐笔数据及快照行情数据
附件提供了某日股票原始快照数据和原始成交数据,用于演示实时生成因子到模拟仿真交易的过程。
//导入模块
clearCachedModules()
use SnapTradeStream::CreateTable
use SnapTradeStream::utils
use SnapTradeStream::OHLCFactors
use SnapTradeStream::StreamFrame
// 导入数据
// 路径需修改
colNames, colTypes = getTradeTickColDefs()
tradeData = select * from loadText("/data/trade.csv", ,table(colNames, colTypes)) order by tradeTime
colNames, colTypes = getSnapshotColDefs()
snapshotData = select * from loadText("/data/snapshot.csv", ,table(colNames, colTypes)) order by tradeTime
6.3.2 建立流表
导入某日成交逐笔数据及快照行情数据之后,需要建立流表,用于后续数据的实时传入和流计算的订阅。
tradeStreamName = "tradeStream"
unsubAndDropAll(tradeStreamName)
sche = schema(tradeData).colDefs
colNames, colTypes = sche.name, sche.typeString
try{
share(streamTable(100:0, colNames, colTypes), tradeStreamName)
}catch(ex){print(ex)}
snapStreamName = "snapshotStream"
unsubAndDropAll(snapStreamName)
sche = schema(snapshotData).colDefs
colNames, colTypes = sche.name, sche.typeString
try{
share(streamTable(100:0, colNames, colTypes), snapStreamName)
}catch(ex){print(ex)}
6.3.3 因子处理函数等参数配置
为适应不同数据源的分钟频特征提取需求,系统设计了统一的处理函数映射字典factorFuncDict与并行度配置字典parallelDict,实现高度模块化和灵活调度。
factorFuncDict = {
"tradeTSE": getTrade1MTSFactor, // 获取成交数据时序聚合引擎指标的函数
"tradeRSE": getTrade1MReactiveFactor, // 获取成交数据响应式状态引擎指标的函数
"snapTSE": getSnap1MTSFactor, // 获取快照数据时序聚合引擎指标的函数
"snapRSE": getSnap1MReactiveFactor, // 获取快照数据响应式状态引擎指标的函数
"mergedRSE": minDerivedKlineFactors // 获取衍生数据响应式状态引擎指标的函数
}
// 定义并行度参数字典
parallelDict = {
"tradeTSE": 5,
"tradeRSE": 5,
"snapTSE": 5,
"snapRSE": 5,
"mergedRSE": 5
}
// 定义模型路径
extraParams = {
"torchModelPath": "/data/pytorch_model.pt",
"lgbmModelPath": "/data/lightgbm_model.txt",
"lowLatency": false
}
-
parallelDict配置了每类因子计算任务的并行线程数。当前每类任务均设为5 个并行任务,可根据实际节点资源和流量进行动态调整。
-
factorFuncDict为一个字典,键为字符串,值为函数,用于指定获取每个流引擎对应的指标及字段定义,该函数默认定义在
klineFactors模块中,比如getTrade1MTSFactor函数,函数定义为:
// 获取成交数据时序聚合引擎指标及字段定义的函数
def getTrade1MTSFactor(){
d = dict(STRING, ANY, true)
d["getAmounts(tradeBSFlag, tradePrice, tradeQty, `B)"] = ["buyAmounts", "DOUBLE"]
d["getAmounts(tradeBSFlag, tradePrice, tradeQty, `S)"] = ["sellAmounts", "DOUBLE"]
colNames2 = ["buySmallAmount", "buyBigAmount", "sellSmallAmount", "sellBigAmount"]
colTypes2 = take("DOUBLE", size(colNames2))
d["calCapitalFlow(buyApplSeqNum,offerApplSeqNum, tradeBSFlag, tradeQty, tradePrice)"] = [colNames2, colTypes2]
d["lastTradeBias(tradeTime, tradeBSFlag, tradePrice, tradeQty)"] = ["tradeBias", "DOUBLE"]
return d
}
该函数返回一个字典,其中键为函数调用字符串,值为包含两个元素的列表:第一个元素为指标名称,第二个元素为字段类型。这种设计使得用户可以根据实际业务需求灵活调整配置。当因子计算逻辑发生变化时,用户只需在此处进行相应修改,无需同步更新流数据表和流计算引擎等底层组件,从而显著降低了因子维护和更新的成本,提升了系统的可扩展性和开发效率。
6.3.4 构建全流程计算框架
定义所需参数后,调用主函数buildStreamFrame进行订阅计算实时生成分钟频因子。
factorDbname = "dfs://factor"
factorTbname = "minDerivedFactorsTb"
buildStreamFrame(factorDbname, factorTbname, tradeStreamName, snapStreamName, factorFuncDict, parallelDict, extraParams)
其中,buildStreamFrame需要传入七个参数,分别为:
| 参数 | 含义 |
|---|---|
| factorDbname | 存储因子的数据库名 |
| factorTbname | 存储因子的数据表名 |
| tradeStreamName | 逐笔成交流表名 |
| snapStreamName | 快照行情流表名 |
| factorFuncDict | 因子处理函数字典 |
| parallelDict | 并行任务个数字典 |
| extraParams | 额外参数,包含模型路径和低延时配置 |
6.3.5 插入数据
将历史数据通过回放功能插入流表中,以此模拟实时流数据流入的效果,从而触发流计算全流程。
replay(
[tradeData, snapshotData],
[objByName(tradeStreamName, true), objByName(snapStreamName, true)],
`tradeTime,
`tradeTime,
100,
false
)
6.3.6 启动推理与模拟撮合
buildStreamFrame负责创建分钟频因子计算管道。回放完成后,等待衍生因子表产生数据,再启动模型推理和模拟撮合:
maxWaitMs = 60000
waitedMs = 0
do {
sleep(500)
waitedMs += 500
} while (size(objByName(factorTbname, true)) == 0 && waitedMs < maxWaitMs)
use SnapTradeStream::streamUtils
derived2Simulator(
factorTbname,
extraParams["lgbmModelPath"],
extraParams["torchModelPath"],
snapStreamName,
false
)
6.3.7 结果展示
实时模拟仿真成交表
以上为实时模拟仿真成交明细结果中的关键字段,其中表字段含义分别为:
| 字段名 | 字段含义 |
|---|---|
| OrderID | 成交的用户订单 ID |
| Symbol | 股票代号 |
| Direction |
买卖方向:
|
| LimitPrice | 委托价格 |
| VolumeTotalOriginal | 委托量 |
| TradeTime | 撮合成交时间 |
| TradePrice | 成交价格 |
| VolumeTraded | 成交量 |
| OrderStatus |
订单状态:
|
| effectiveTime | 订单开始时间 |
| expireTime | 订单结束时间 |
| cumTradeQty | 累计成交量 |
| tradeAvgPrice | 平均成交价格 |
6.3.8 清理环境
通过执行以下代码清理流计算环境:
clearStreamEnv()
each(undef{, SHARED}, [tradeStreamName, snapStreamName])
6.4 实时低延时计算
DolphinDB 从3.00.4版本开始,支持在时序聚合引擎和响应式引擎指定低延时参数lowLatency,设置为true时,引擎将采用行式低延时算法,适用于对低延要求极高、且每次写入数据量较小的场景,本文仅介绍在 dolphindb 下的场景,若追求更低时延,可以使用swordfish的低延时计算功能。以创建时序聚合引擎为例,具体使用方式为:
tse = createTimeSeriesEngine(
name= tseName,
windowSize=60000,
step=60000,
metrics=tseMetrics,
dummyTable=rawStream,
outputTable=rse,
timeColumn=`tradeTime,
useSystemTime=false,
keyColumn=`securityID,
useWindowStartTime=false,
parallelism=parallelTSE,
lowLatency = true, // 添加该参数即可
outputElapsedMicroseconds = true
)
在本文介绍的模块中,可以通过:
extraParams = {
"torchModelPath": "/data/pytorch_model.pt",
"lgbmModelPath": "/data/lightgbm_model.txt",
"lowLatency": true
}
支持在整个流计算框架中的所有响应式状态引擎和时序聚合引擎都采用低延时模式。本部分模拟了实时计算的场景,测试在关闭和开启低延时参数的情况下,不同数据量的流计算引擎计算耗时情况。
测试结果
选取全天 100、200、500、1000 支股票的真实快照数据和成交数据,在关闭与开启低时延参数两种情况下,测试 TSE(流数据日级时间序列引擎:createDailyTimeSeriesEngine)和 RSE(响应式状态引擎: createReactiveStateEngine)的平均计算耗时(单位:微秒)。括号中指标表示开启低时延参数后的耗时与关闭该参数时耗时的比值。
| SampleCounts | Snapshot | Trade | DerivedFactors | |||||||
|---|---|---|---|---|---|---|---|---|---|---|
| TSE | RSE | TSE | RSE | RSE | ||||||
| FALSE | TRUE | FALSE | TRUE | FALSE | TRUE | FALSE | TRUE | FALSE | TRUE | |
| 100 | 59.13 | 39.28 (0.66) | 11.72 | 9.97 (0.85) | 60.67 | 53.01 (0.87) | 2.11 | 0.79 (0.37) | 151.51 | 57.48 (0.38) |
| 200 | 61.46 | 40.11 (0.65) | 12.09 | 9.66 (0.80) | 62.46 | 59.80 (0.96) | 2.07 | 0.93 (0.45) | 185.12 | 59.24 (0.32) |
| 500 | 64.47 | 42.27 (0.66) | 12.46 | 9.87 (0.79) | 63.83 | 57.75 (0.90) | 2.16 | 1.05 (0.49) | 317.54 | 72.69 (0.23) |
| 1000 | 66.13 | 44.06 (0.67) | 12.29 | 9.84 (0.80) | 67.25 | 57.76 (0.86) | 2.19 | 1.13 (0.52) | 1,587.00 | 94.72 (0.06) |
可以观察到,当开启低时延参数时,所有流计算引擎的计算耗时均明显小于关闭该参数的表现。
7. Orca 计算平台
面对大规模集群部署、复杂任务依赖和高可用性等需求,手动管理多个引擎与订阅关系日益复杂。DolphinDB 推出的 Orca 实时计算平台采用声明式 API 与集中式智能调度,简化了流处理逻辑的设计、部署和运维全流程,提升了系统稳定性与资源利用效率。SnapTradeStream 模块基于 Orca 构建了与第六章功能相同的实时因子计算与仿真交易流程,代码更为清晰系统化。
下图为 Orca 流程的整体架构:
如图所示,Orca 实时计算平台利用集群计算资源,基于成交和快照原始数据实时生成分钟频衍生因子。执行节点订阅因子数据进行实时推理并输出交易信号,交易信号用于构造委托订单,发送至模拟撮合订单管理引擎,最终输出成交明细。
7.1 Orca 计算平台应用实践
传统流计算需手动完成数据接入、引擎配置、因子函数注册和模型加载等环节,过程繁琐。Orca 平台通过统一流图结构,将因子计算封装为声明式流程,模型推理与模拟仿真交易在本节点执行,提升了开发效率。本章基于 Orca 实现上述系统的搭建。
7.1.1 定义参数
以下代码完成 Orca 分钟因子框架的初始化:先导入成交与快照样本数据,再建同结构流表,并配置因子函数、并行度、模型路径与结果落库位置。
use SnapTradeStream::CreateTable
use SnapTradeStream::utils
use SnapTradeStream::KlineFactors
use SnapTradeStream::OrcaFrame
// 导入数据
colNames, colTypes = getTradeTickColDefs()
tradeData = select * from loadText("/data/trade.csv", ,table(colNames, colTypes)) order by tradeTime
colNames, colTypes = getSnapshotColDefs()
snapshotData = select * from loadText("/data/snapshot.csv", ,table(colNames, colTypes)) order by tradeTime
// 创建流表
tradeStreamName = "tradeStreamOrca"
unsubAndDropAll(tradeStreamName)
sche = schema(tradeData).colDefs
colNames, colTypes = sche.name, sche.typeString
try{
share(streamTable(100:0, colNames, colTypes), tradeStreamName)
}catch(ex){print(ex)}
snapStreamName = "snapshotStreamOrca"
unsubAndDropAll(snapStreamName)
sche = schema(snapshotData).colDefs
colNames, colTypes = sche.name, sche.typeString
try{
share(streamTable(100:0, colNames, colTypes), snapStreamName)
}catch(ex){print(ex)}
// 定义因子函数字典
factorFuncDict = {
"tradeTSE": getTrade1MTSFactor, // 获取成交数据时序聚合引擎指标的函数
"tradeRSE": getTrade1MReactiveFactor, // 获取成交数据响应式状态引擎指标的函数
"snapTSE": getSnap1MTSFactor, // 获取快照数据时序聚合引擎指标的函数
"snapRSE": getSnap1MReactiveFactor, // 获取快照数据响应式状态引擎指标的函数
"mergedRSE": minDerivedKlineFactors // 获取衍生数据响应式状态引擎指标的函数
}
// 定义并行度参数字典
parallelDict = {
"tradeTSE": 5,
"tradeRSE": 5,
"snapTSE": 5,
"snapRSE": 5,
"mergedRSE": 5
}
// 定义模型路径
extraParams = {
"torchModelPath": "/data/pytorch_model.pt",
"lgbmModelPath": "/data/lightgbm_model.txt"
}
// 定义持久化库表路径
factorDbname = "dfs://factor"
factorTbname = "minDerivedFactorsTbOrca"
| 参数 | 含义 | 说明 |
|---|---|---|
| factorDbname | 存储因子的数据库名 | 示例为dfs://factor |
| factorTbname | 存储因子的数据表名 | 示例为minDerivedFactorsTbOrca |
| tradeStreamName | 逐笔成交流表名 | 示例为tradeStreamOrca,需与导入的成交数据同结构 |
| snapStreamName | 快照行情流表名 | 示例为snapshotStreamOrca,需与导入的快照数据同结构 |
| factorFuncDict | 因子处理函数字典 | 键为不同名称,值为对应因子定义函数,详见下表 |
| parallelDict | 并行任务个数字典 | 键与factorFuncDict对齐,值为各阶段并行任务数 |
| extraParams | 额外参数,包含模型路径和低延时配置 | 传入推理模型路径等扩展配置,详见下表 |
factorFuncDict各键含义如下:
| 参数 | 含义 | 说明 |
|---|---|---|
| tradeTSE | 成交时序聚合因子函数 | 对应getTrade1MTSFactor,计算成交数据 1 分钟时序聚合指标 |
| tradeRSE | 成交响应式状态因子函数 | 对应getTrade1MReactiveFactor,计算成交数据响应式状态指标 |
| snapTSE | 快照时序聚合因子函数 | 对应getSnap1MTSFactor,计算快照数据 1 分钟时序聚合指标 |
| snapRSE | 快照响应式状态因子函数 | 对应getSnap1MReactiveFactor,计算快照数据响应式状态指标 |
| mergedRSE | 衍生因子函数 | 对应minDerivedKlineFactors,在合并结果上计算衍生 K 线因子 |
parallelDict各键含义如下:
| 参数 | 含义 | 说明 |
|---|---|---|
| tradeTSE | 成交时序聚合并行度 | 示例为5,按标的分组并发计算成交 1 分钟时序聚合因子 |
| tradeRSE | 成交响应式状态并行度 | 示例为5,按标的分组并发计算成交响应式状态因子 |
| snapTSE | 快照时序聚合并行度 | 示例为5,按标的分组并发计算快照 1 分钟时序聚合因子 |
| snapRSE | 快照响应式状态并行度 | 示例为5,按标的分组并发计算快照响应式状态因子 |
| mergedRSE | 衍生因子并行度 | 示例为5,按标的分组并发计算衍生 K 线因子 |
extraParams各键含义如下:
| 参数 | 含义 | 说明 |
|---|---|---|
| torchModelPath | PyTorch 模型路径 | 示例为/data/pytorch_model.pt,供下游推理加载 |
| lgbmModelPath | LightGBM 模型路径 | 示例为/data/lightgbm_model.txt,供下游推理加载 |
7.1.2 创建框架
通过调用模块中的buildOrcaFrame函数搭建 Orca 流计算框架,参数介绍参考 6.3.4 小节,该函数主要做以下流程操作:
-
清理现有环境,包括流图、流表、流引擎及相关订阅,确保无残留数据干扰。
-
通过 Orca 接口创建流计算框架,实现基于原始快照数据和成交数据的分钟频衍生因子指标生成。
-
订阅这些分钟频衍生指标数据,利用模型进行实时推理,生成交易信号,随后构造委托订单并插入模拟撮合订单管理引擎,最终实现实时模拟交易结果的生成。
buildOrcaFrame(factorDbname, factorTbname, tradeStreamName, snapStreamName, factorFuncDict, parallelDict, extraParams)
7.1.3 回放数据
通过回放历史数据,触发 Orca 流计算与模拟仿真交易全流程,实时输出模拟交易结果。
replay(
[tradeData, snapshotData],
[objByName(tradeStreamName, true), objByName(snapStreamName, true)],
`tradeTime,
`tradeTime,
100,
false
)
7.1.4 清理环境
通过执行以下代码清理 orca 环境
clearOrcaEnv()
each(undef{, SHARED}, [tradeStreamName, snapStreamName])
7.2 Orca 优势总结
以计算原始成交数据的分钟频指标为例,使用 Orca 替代传统流计算流程进行因子计算。整个流程从原始数据输入,经过响应式状态引擎、时序聚合引擎,最终写入持久化表。相比传统流计算,Orca 仅需通过链式函数明确声明所需步骤,无需繁琐定义表结构、共享表等与业务无关的配置,极大简化开发流程,提高开发效率和代码可读性:
tradePipeline = g.source(tradeStreamName, 1:0, colNames1, colTypes1)
.dailyTimeSeriesEngine(
windowSize = 60000,
step = 60000,
metrics = tradetseMetrics,
timeColumn = `tradeTime,
useSystemTime=false,
keyColumn = `securityID,
useWindowStartTime=false,
parallelism=tradeparallelTSE,
sessionBegin=[09:30:00.000, 12:59:00.000],
sessionEnd=[11:30:00.000, 14:57:00.000]
)
.reactiveStateEngine(
metrics = traderseMetrics,
keyColumn = `securityId,
parallelism = tradeparallelRSE
)
.sink(factorDbname+"/"+tradeMinStreamTbname)
总结而言,Orca 流计算平台通过声明式 API 和集中式智能调度,系统性地替代了传统手工构建的分布式因子计算流程。该平台统一并简化了从原始行情数据接入、分钟频聚合到复杂因子计算的全链路数据处理过程,显著提升了开发效率,降低了运维复杂度,增强了系统可靠性,同时保障了资源利用率和业务连续性。凭借这些优势,Orca 在量化场景下实现了实时因子的高效、稳定计算,成为强有力的平台级基础支撑。
8. 总结
本教程系统讲解了如何在 DolphinDB 中利用股票 Level 2 快照数据与成交行情数据,实时生成分钟频因子及衍生因子,并基于此进行推理计算。详细演示了分别运用传统流计算低延时框架和Orca 框架进行实时分钟频因子及衍生因子计算、推理和模拟仿真交易的完整解决方案。通过对比两种方法的实现流程和特点,帮助用户更直观地理解各方案的适用场景与优势,进而为不同需求提供多样化的实时计算方案。
需要注意的是,本教程中所介绍的分钟频因子和衍生因子的生成规则,可能与用户的实际应用场景有所不同。用户可以根据自身需求进行调整,从而快速完成项目开发。
附件
版本及要求
-
本教程代码基于 DolphinDB3.00.4开发,建议用户使用3.00.4及以上版本 。
-
本教程用到多个插件,请预先安装与加载:
-
lgbm -
LibTorch -
MatchingEngineSimulator -
OrderManagementEngine
-
快照原始数据表结构
| 字段名称 | 数据类型 | 数据说明 |
|---|---|---|
| Market | SYMBOL | 市场代码 |
| tradeTime | TIMESTAMP | 交易时间 |
| MDStreamID | INT | 行情类别 |
| securityID | SYMBOL | 证券代码 |
| numImageStatus | INT | 快照图像状态 |
| securityIDSource | INT | 证券代码来源 |
| tradingPhaseCode | SYMBOL | 交易阶段 |
| preCloPrice | DOUBLE | 前收盘价 |
| numTrades | INT | 成交笔数 |
| totalVolumeTrade | LONG | 总成交量 |
| totalValueTrade | DOUBLE | 总成交额 |
| lastPrice | DOUBLE | 最新成交价 |
| openPrice | DOUBLE | 开盘价 |
| highPrice | DOUBLE | 最高价 |
| lowPrice | DOUBLE | 最低价 |
| closePrice | DOUBLE | 收盘价 |
| difPrice1 | DOUBLE | 涨跌一 |
| difPrice2 | DOUBLE | 涨跌二 |
| PE1 | DOUBLE | 市盈率一 |
| PE2 | DOUBLE | 市盈率二 |
| preCloseIOPV | DOUBLE | 前收盘 IOPV |
| IOPV | DOUBLE | IOPV净值 |
| totalBuyQty | LONG | 总买入委托量 |
| weightedAvgBuyPx | DOUBLE | 加权买一价 |
| totalOfferQty | LONG | 总卖出委托量 |
| weightedAvgOfferPx | DOUBLE | 加权卖一价 |
| upLimitPx | DOUBLE | 涨停价 |
| downLimitPx | DOUBLE | 跌停价 |
| openInt | LONG | 持仓量 |
| optPremiumRatio | DOUBLE | 期权权利金比例 |
| offerPrice | DOUBLE[] | 卖一~卖十价 |
| buyPrice | DOUBLE[] | 买一~买十价 |
| offerOrderQty | LONG[] | 卖一~卖十量 |
| buyOrderQty | LONG[] | 买一~买十量 |
| buyNumOrders | LONG[] | 买一~买十委托笔数 |
| offerNumOrders | LONG[] | 卖一~卖十委托笔数 |
| ETFBuyNumber | INT | ETF 申购笔数 |
| ETFBuyAmount | LONG | ETF 申购份额 |
| ETFBuyMoney | DOUBLE | ETF 申购金额 |
| ETFSellNumber | INT | ETF 赎回笔数 |
| ETFSellAmount | LONG | ETF 赎回份额 |
| ETFSellMoney | DOUBLE | ETF 赎回金额 |
| yieldToMatu | DOUBLE | 到期收益率 |
| withdrawBuyNumber | INT | 买入撤单笔数 |
| withdrawBuyAmount | LONG | 买入撤单量 |
| withdrawBuyMoney | DOUBLE | 买入撤单金额 |
| withdrawSellNumber | INT | 卖出撤单笔数 |
| withdrawSellAmount | LONG | 卖出撤单量 |
| withdrawSellMoney | DOUBLE | 卖出撤单金额 |
| totalBuyNumber | INT | 买单总笔数 |
| totalOfferNumber | INT | 卖单总笔数 |
| maxBuyDur | INT | 最长买单挂单时长 |
| maxSellDur | INT | 最长卖单挂单时长 |
| buyNum | INT | 买入档数 |
| sellNum | INT | 卖出档数 |
| localTime | TIME | 本地时间戳(采集机时间) |
| seqNo | INT | 本地序列号 |
| offerOrders | LONG[] | 卖出委托队列 |
| buyOrders | LONG[] | 买入委托队列 |
成交原始数据表结构
| 字段名称 | 数据类型 | 数据说明 |
|---|---|---|
| channelNo | INT | 频道代码 |
| applSeqNum | LONG | 消息序号 |
| MDStreamID | INT | 行情类别 |
| buyApplSeqNum | INT | 买方委托序号 |
| offerApplSeqNum | INT | 卖方委托序号 |
| securityID | SYMBOL | 证券代码 |
| securityIDSource | INT | 证券代码源 |
| tradePrice | DOUBLE | 成交价格(元) |
| tradeQty | LONG | 成交数量(股) |
| execType | INT | 成交类型 |
| tradeTime | TIMESTAMP | 成交时间 |
| localTime | TIME | 本地采集时间 |
| seqNo | INT | 本地序号 |
| tradeBSFlag | SYMBOL | 主动买卖标志(B-主动买,S-主动卖) |
| market | SYMBOL | 市场代码 |
| dataStatus | INT | 数据状态 |
| tradeIndex | INT | 成交索引 |
| tradeMoney | DOUBLE | 成交金额(元) |
| bizIndex | LONG | 业务索引 |
模块脚本与数据
-
流计算模块:SnapTradeStream.zip
-
调用脚本:
-
传统流计算:stream.dos
-
orca :orca.dos
-
-
训练模型 demo:
-
LightGBM 模型:lightgbm_model.txt
-
PyTorch 模型:pytorch_model.pt
-
-
快照数据:snapshot.csv
-
成交数据:trade.csv
-
OrderManagementEngine 插件:OrderManagementEngine.zip
