委托订单还原引擎使用最佳实践

在基于 Level-2 行情开发高频策略时,逐笔委托(Order)和逐笔成交(Trade)是最基础的数据来源。

很多开发者在将同一套策略同时应用于沪深两市时,都会遇到一个容易忽略的问题:基于原始的逐笔委托数据统计委托量或订单流特征时,同样的统计逻辑,在深交所可以直接计算得到结果,但在上交所直接计算得到的结果往往会低于实际值。

例如,统计某只股票最近 60 秒新增的买单委托量:

过去60秒新增买单委托量 = 所有新到达买单委托数量之和

对于深交所,可以直接在逐笔委托数据上统计。

而对于上交所,即使统计逻辑完全一致,结果也会偏小。

造成这一差异的原因,并不是计算逻辑错误,而是两个交易所逐笔数据的发布规则不同

因此,在许多高频策略中,开发者需要:

  • 为沪深两市维护两套数据预处理逻辑;

  • 在上交所数据中额外关联逐笔成交,反推缺失的委托;

  • 对处理后的数据进行统一口径处理,确保沪深两市数据具有可比性。

这种额外的数据处理工作不仅增加了开发复杂度,也提高了策略迁移和维护成本。

为了解决这一问题,DolphinDB 提供了上交所委托订单还原引擎( createOrderReconstituteEngine ,能够恢复上交所缺失的原始委托记录,使输出数据与深交所逐笔委托具有一致的数据语义。

注意:本文涉及的完整代码见附录,实例代码建议在 3.00.4 及以上版本的 DolphinDB server 上运行。

1. 沪深交易所逐笔数据差异

本节介绍沪深交易所逐笔数据在委托记录上的主要差异以及这些差异对后续计算的影响。

1.1 深交所完整发布

深交所的逐笔委托数据能够完整记录进入撮合系统的委托。无论委托最终是否成交,逐笔委托表都会记录对应的原始委托信息。

因此,基于逐笔委托数据即可直接统计委托量、订单流等特征,无需额外关联逐笔成交数据来补全委托信息。

1.2 上交所不完整发布

上交所对能够立即撮合的委托做了精简,具体分为以下两种情形。

情形一:委托提交后立即全部成交。此时上交所不会再发送原始委托记录,只保留对应的成交记录。

  • 如图 1-1 所示,一笔委托量为 3,700 的买单(委托号 218,629)提交后,立即与市场上的 4 笔卖单(委托号分别为 204,322、34,204、137,957、210,374)完全成交。逐笔数据中会出现连续的 4 条成交记录(1300 + 1500 + 500 + 400 = 3700),但找不到委托量为 3,700 的原始委托记录。

1. 图 1-1 上交所即时全部成交情形

情形二:委托提交后立即部分成交。此时逐笔数据中会先出现成交记录,再出现一条委托记录,且委托记录中的委托量是剩余未成交量,而非原始委托量;原始委托量等于成交量与剩余委托量之和。

  • 如图 1-2 所示,一笔委托量为 2,000 的卖单(委托号 218,630)提交后,其中 531 股立即与市场上一笔买单(委托号 214,265)成交,剩余 1,469 股作为一条新的委托记录发送(该委托后续在 09:30:17 被撤单)。逐笔数据中同样找不到委托量为 2,000(531 + 1469)的原始委托记录。

2. 图 1-2 上交所即时部分成交情形

1.3 数据差异的影响

很多高频指标需要统计过去一段时间市场新增的委托数量例如:

  • 最近 60 秒新增买单数量

  • 最近 60 秒新增卖单数量

  • 最近 5 分钟委托量波动率

这些指标关注的是订单到达情况,而不是最终是否成交。

如果逐笔委托数据中缺少立即成交的委托记录,基于原始数据直接统计就会低估实际的委托流入量。

因此,对于上交所数据,需要首先恢复这些缺失的委托,再进行统计分析。

2. 委托订单还原引擎介绍

委托订单还原引擎的作用,就是根据上交所逐笔成交与逐笔委托之间的关联关系,实时还原缺失的原始委托信息。

2.1 输出数据结构

输出表在输入表的基础上新增两个字段。具体示例见图 2-1,新增字段位于表格的最后两列。

表 2-1 新增字段说明表

字段 类型 说明
orderMark INT 记录逐笔数据的来源
orderIndex LONG 记录输出顺序,从 0 开始递增。

其中,orderMark 用于区分不同来源的数据。

表 2-2 orderMark来源表

orderMark 委托记录 成交记录
0 原始委托 原始成交
1 恢复的全部即时成交委托 全部即时成交
2 恢复的部分即时成交委托 部分即时成交

2.2 委托还原引擎工作原理

情形一:对于立刻且全部成交的委托单,会补上一条原始委托信息。

如图 2-1 所示,相比于原始数据,还原后的数据新增了一条委托量为 3,700(3700 = 1300 + 1500 + 500 + 400)的原始委托数据。

3. 图 2-1 即时全部成交情形还原后的委托示例

注:

  • 最后一列 orderIndex 用于替代原来的数据中的 ApplSeqNum,作为通道内的单增序号列。ApplSeqNum 从 1 开始递增,而 orderIndex 从 0 开始递增,如图 2-2 所示。

4. 图 2-2 ApplSeqNum 和 orderIndex 字段编号
  • 还原具体步骤:

    • 引擎读到第一笔成交单(图 2-1 第二行 ApplSeqNum = 243,059),发现它是由一笔没有记录的委托单触发的。引擎通过 BuyNo 找到该委托对应的所有成交单,用这些成交单的成交量汇总还原出原始委托量。

    • 引擎还原一条委托量为 3,700 的原始委托,并将其插入当前 orderIndex 所对应的位置(即 243,058)。为了尽可能保留溯源信息,它的 ApplSeqNum 继承了触发它的第一笔成交单的值(243,059)。

    • 第一笔成交单的 orderIndex 顺延为 243,059,后续记录的 orderIndex 依次顺延。

  • 由于是即时全部成交,还原的委托和原始的成交数据 orderMark 都为 1。

情形二:对于立刻部分成交的委托单,会补上一条原始委托信息,原来的剩余委托记录不再保留,因为它已经包含在恢复后的原始委托之中。

如图 2-3 所示,相比于原始数据,增加了一条委托量为 2,000(2000 = 531 + 1469)的原始委托数据,并减少了一条委托量为 1469 的部分委托数据。

5. 图 2-3 即时部分成交情形还原后的委托示例

注:

  • 当一笔 2,000 股的卖单下达,并立刻与盘口的 531 股成交时,上交所原始数据流顺序是先推成交,再推剩余的委托,如图 2-4 所示:

    6. 图 2-4 部分成交情形的成交和委托消息推送顺序
    1. 先收到(ApplSeqNum = 243,063):一笔 531 股的逐笔成交。

    2. 后收到 (ApplSeqNum = 243,064):一笔 1,469 股的逐笔委托(这是扣除成交后的剩余委托)。

  • 还原具体步骤:

    • 引擎将后收到的 1,469 股剩余委托,加上已知的 531 股成交,反推出了一笔 2,000 股的原始委托单。这条还原出的数据,保留了它原本作为剩余委托时的身份标识,即 ApplSeqNum = 243,064。

    • 引擎把较小的逻辑序号 orderIndex = 243,063 分配给还原的 2,000 股委托记录,使该委托记录在成交记录之前。

    • 引擎把较大的逻辑序号 orderIndex = 243,064 分配给 531 股成交记录。该成交记录原始的 ApplSeqNum 为 243,063,通过调整 orderIndex,确保还原的委托记录先于该成交记录。

3. 输入数据准备

委托订单还原引擎的计算正确性依赖于输入逐笔数据的完整性、枚举值规范性以及同一通道内的事件顺序。使用引擎前,需要先将逐笔委托和逐笔成交标准化为统一表结构,确保该表能够准确表达委托、成交、撤单以及通道内顺序信息。

本节讲解对输入数据的要求和具体示例。

3.1 数据处理注意事项

  • 输入数据:成交和委托的合并表

委托订单还原引擎的输入为一张同时包含逐笔成交和逐笔委托的数据表。

将逐笔成交、逐笔委托处理成相同的表结构可以存储到同一张输入表中。合并时新增 SourceType 字段区分消息来源:0 表示逐笔委托,1 表示逐笔成交。后续引擎会根据成交与委托之间的委托号关系,识别全部即时成交、部分即时成交等场景,并输出补全后的委托流。

  • 完整的输入字段【必填】

输入数据至少需要包含以下 8 个字段。输入表中字段名可以自定义,只需要在创建委托还原引擎时需要通过 inputColMap 参数将引擎所需的逻辑字段映射到实际列名。

表 3-1 必选输入字段表

inputColMap 的 key inputColMap 的 value 对应输入表中列类型 字段含义说明 inputColMap 的 value 示例
codeColumn SYMBOL 标的代码 SecurityID
typeColumn INT

交易类型:

  • 如果是逐笔委托单,则:1 表示市价;2 表示限价;3 表示本方最优;10 表示撤单(仅上交所);11 表示市场状态(仅上交所)

  • 如果是逐笔成交单,则:0 表示成交;1 表示撤单(仅深交所)

Type
priceColumn LONG

价格,整型

比如真实价格是浮点数,需要保留 4 位小数,则真实价格 × 10000 后取整

Price
qtyColumn LONG 数量(股数) Qty
buyOrderColumn LONG
  • 逐笔成交:对应其原始成交中的买方委托序号。

  • 逐笔委托:

    • 上交所:填充原始委托中的原始订单号,即上交所在新增、删除订单时用以标识订单的唯一编号

    • 深交所:填充 0。此字段为深交所为了补全上交所数据格式而增加的冗余列

BuyNo
sellOrderColumn LONG
  • 逐笔成交:对应其原始成交中的卖方委托序号。

  • 逐笔委托:

    • 上交所:填充原始委托中的原始订单号,即上交所在新增、删除订单时用以标识订单的唯一编号

    • 深交所:填充 0。此字段深交所为了补全上交所数据格式而增加的冗余列

SellNo
sideColumn INT

买卖方向:1 表示买单;2 表示卖单

说明:

  • 委托单的 BSFlag,必填

  • 撤单的 BSFlag 由原始委托单决定买卖方向,必填

  • 成交单的 BSFlag,不影响结果,非必填

BSFlag
msgTypeColumn INT

数据类型:

  • 0 表示逐笔委托;

  • 1 表示逐笔成交;

  • -1 表示产品状态。

SourceType
  • 推荐保留字段【选填】

除了上述必须包含的 8 个字段外,输入表中可以冗余其他字段,并在输出表中一同输出。

比如 DolphinDB 中已有快照合成引擎,且快照合成引擎也是基于逐笔合并数据进行处理。上交所还原引擎的输出可以作为快照合成引擎的输入。

表 3-2 推荐保留字段表

示例列名 列类型 字段含义说明 推荐保留理由
TradeDate DATE 交易日期 保留原始数据的交易日期
Time TIME 交易时间

保留原始数据的交易时间

是快照合成引擎的必填列,对应 createOrderBookSnapshotEngine 函数中 inputColMap 参数的 timeColumn

ApplSeqNum LONG 一个通道内从 1 开始递增的逐笔数据序号。深交所为 appseqlnum 字段;上交所为 bizIndex 字段。

是一个通道内逐笔数据的唯一标识。

是快照合成引擎的必填列,对应 createOrderBookSnapshotEngine 函数中 inputColMap 参数的 seqColumn

ChannelNo INT 行情通道号

逐笔数据的唯一键:ChannelNo + ApplSeqNum;

因此,基于逐笔数据的处理(比如上交所委托还原引擎和快照合成引擎)都需要分通道;

一个通道一个引擎。

ReceiveTime NANOTIMESTAMP 逐笔数据的接收时间 是快照合成引擎的必填列,对应 createOrderBookSnapshotEngine 函数中 inputColMap 参数的 receiveTime
  • 规范的枚举值

上交所委托还原引擎对相关字段的枚举值有明确约定。输入表中的枚举值必须符合约定才能得到正确结果。比如 BSFlag 字段,约定 1 代表买方向, 2 代表卖方向。

  • 单个引擎处理范围:对于同一个引擎,至多只能输入一天一个通道的全部标的

委托订单还原依赖成交与委托之间的委托号关系、通道内逐笔消息顺序,以及引擎维护的通道内状态。因此,输入同一个引擎的数据应来自同一交易日、同一通道号,并按通道内消息顺序写入。

时间范围:单个引擎实例只能处理一天的数据,不能跨交易日处理。如果需要处理多天数据,则需要每天重新创建引擎实例。

标的范围:单个引擎实例只能处理一个通道的数据。如果需要处理多个通道,则需要为每个通道分别创建引擎实例。

输入数据顺序:单通道内按 ApplSeqNum 按升序排序。

3.2 历史批计算数据处理

本节提供批量合并上交所逐笔成交和逐笔委托至同一张表中的示例。

输入数据

本节提供测试数据压缩包:SH20260429_db.zip,其中包含 2026.04.29 当天 2 支上交所股票的完整逐笔委托数据(SHorder)和逐笔成交数据(SHtrans)。

将数据压缩包上传至 DolphinDB 部署服务器并解压缩:

  • 如果将数据文件夹放置在 DolphinDB 节点主目录下,可以直接运行示例脚本。可在 DolphinDB 客户端执行 getHomeDir() 查询节点主目录。例如 /home/user/DolphinDB/server/SH20260429_db

  • 如果将数据放置在其他自定义目录,则需要修改示例脚本中的 sh_db 变量,将其设置为实际的数据目录后再运行脚本。

示例脚本

// 读取指定的磁盘目录下的数据,使用 getHomeDir() 拼接得到绝对路径
sh_db = getHomeDir() + "/SH20260429_db"

// 加载逐笔委托和逐笔成交表
SHorder = loadTable(sh_db, "SHorder")
SHtrans = loadTable(sh_db, "SHtrans")

// 创建合并表 SHOrderTrans
name = ["TradeDate","SecurityID","Time","SourceType","Type","Price","Qty",
        "BSFlag","BuyNo","SellNo","ApplSeqNum","ChannelNo","Market","ReceiveTime"]
type = ["DATE","SYMBOL","TIME","INT","INT","LONG","LONG",
        "INT","LONG","LONG","LONG","INT","SYMBOL","TIMESTAMP"]
SHOrderTrans = table(1:0, name, type)

// 指定枚举值的映射关系
orderTypeMap = dict(["S", "A", "D"], [-1, 2, 10])
BSFlagMap = dict(["B", "S"], [1, 2])

// 逐笔委托数据格式处理: 加载order表并处理为特定的表结构	
orderTemp = 
    select trade_date as TradeDate, 
        security_id as SecurityID, 
        trade_time as Time,
        iif(order_type=="S", -1, 0) as SourceType,  
        orderTypeMap[order_type] as Type,
        long(order_price * 10000) as Price,
        order_qty as Qty,
        BSFlagMap[side] as BSFlag,
        order_no as BuyNo,
        order_no as SellNo,
        appl_seq as ApplSeqNum,
        channel_no as ChannelNo,
        exchange as Market,
        update_time as ReceiveTime
        from SHorder
SHOrderTrans.append!(orderTemp)

// 逐笔成交数据格式处理: 加载transaction表并处理为特定的表结构
tradeTemp = 
    select 
        trade_date as TradeDate, 
        security_id as SecurityID, 
        trade_time as Time,
        1 as SourceType,
        0 as Type,
        long(trade_price * 10000) as Price,
        trade_qty as Qty,
        BSFlagMap[side].nullFill(iif(sell_no > buy_no, 2, 1)) as BSFlag, // side缺失时用买卖单号大小关系推断方向
        buy_no as BuyNo,
        sell_no as SellNo,
        appl_seq as ApplSeqNum,
        channel_no as ChannelNo,
        exchange as Market,
        update_time as ReceiveTime
    from SHtrans
SHOrderTrans.append!(tradeTemp)


// 处理逐笔数据的顺序: 按照 ApplSeqNum 进行排序
SHOrderTrans = select * from SHOrderTrans order by ApplSeqNum
  • 逐笔委托数据处理说明:

    • 逐笔委托数据中包含产品状态信息,所以需要处理数据源字段。当 order_type="S" 时,表示该记录为产品状态,对应的 SourceType 设置为 -1,示例脚本中通过 iif(order_type=="S", -1, 0) 完成该处理。

    • 上交所委托类型通过 orderTypeMap 映射为统一编码:S 表示产品状态,映射为 -1;A 表示新增委托,映射为 2;D 表示撤单,映射为 10。

    • 上交所还原引擎要求输入价格列为整型。所以示例脚本中通过 long(Price*10000) 对价格做保留 4 位小数的处理。

    • 上交所逐笔委托数据中的 order_no 用于标识委托单的唯一编号,逐笔成交数据中的 buy_no 和 sell_no 分别对应买方和卖方委托单编号。原始委托数据中没有 buy_no 和 sell_no 字段,所以处理时使用委托数据中的 order_no 进行填充。

  • 逐笔成交数据处理说明:

    • 上交所逐笔成交数据没有单独的成交类型字段,示例中将成交记录的 Type 统一填充为 0。

    • 上交所还原引擎要求输入价格列为整型。所以示例脚本中通过 long(Price*10000) 对价格做保留 4 位小数的处理。

    • 当成交记录中的 side 缺失时,示例使用 iif(sell_no > buy_no, 2, 1) 根据买卖委托号大小关系推断买卖方向,并通过 nullFill 补齐 BSFlag。

  • 逐笔数据顺序处理说明:

    • 合并结果写入内存表 SHOrderTrans。由于委托还原引擎需要依赖同一通道内的逐笔顺序,后续输入引擎前应按交易日和 ChannelNo 过滤,并按 ApplSeqNum 升序排序。

如图 3-1 所示,逐笔委托与逐笔成交合并后生成 SHOrderTrans 表。

7. 图 3-1 逐笔委托和逐笔成交合并表

注:深交所逐笔委托数据已经完整披露,不需要、也不应写入上交所委托订单还原引擎。该引擎的处理逻辑针对上交所“立即全部成交不发送原始委托、立即部分成交只发送剩余委托”的发布规则设计;如果将深交所数据写入引擎,引擎会按照上交所规则尝试识别和改写委托流,导致输出数据与深交所原始逐笔数据不一致。因此,深交所数据只需按照统一字段口径完成标准化处理,后续可直接接入快照合成或因子计算流程。

3.3 实时流数据处理

在实时数据接入委托还原引擎时,必须保证逐笔成交和逐笔委托以交易所发送的真实顺序有序写入。

因此,实时输入数据处理可以概括为以下流程:

  1. 在 DolphinDB 中创建共享流数据表,用于接收处理好的逐笔数据。表结构要求参考 3.1 章节。

  2. 订阅共享流数据表,将合并后的逐笔数据实时写入委托还原引擎中。

  3. 通过行情 SDK 同时订阅逐笔成交和逐笔委托数据,并将同一个通道的逐笔成交和逐笔委托数据写入同一张共享流数据表。即,在行情 SDK 的回调函数中解析收到的逐笔成交和逐笔委托数据,并完成类型转换、枚举值映射以及字段填充等数据处理工作。数据处理示例参考 3.2 章节。

一般推荐以下方式将实时数据源接入 DolphinDB server 中:

  • 方案一:使用 DolphinDB 行情数据源插件。

    • DolphinDB 的行情数据源插件已经支持将委托和成交接收到同一张表中,并且内置实现了枚举值转换,可以直接接收到委托还原引擎期望的输入数据。

      • amdQuote 插件:amdQuote 订阅时选择 dataType="orderExecution"

      • NSQ 插件:订阅时选择dataType="orderTrade"

      • MDL 插件:订阅时选择 svrID="SHL2_ORDER_AND_TRANSACTION"

      • WindTDF 插件:订阅时选择 dataType="orderTrade"

      • EFH 插件:订阅时选择 dataType="orderExecution"

      • CSM 插件:订阅时选择 dataType="SSEL2_Tick"

      • INSIGHT 插件:订阅时选择 marketDataTypes="MD_ORDER_TRANSACTION"

      • XTP 插件:订阅时在 tableDict 参数里选择接收逐笔合并表(TickByTickTable)

  • 方案二:使用 DolphinDB API 编写外部程序。

    • 自行开发数据接收程序,比如 C++ 程序。

    • 推荐使用 MultithreadedTableWriter 接口往 DolphinDB 的共享流表内实时写入数据。

    • 接收程序内的处理逻辑可以参考上述流程和插件源码。

4. 使用委托订单还原引擎

本节演示如何使用委托订单还原引擎重建上交所逐笔数据。

4.1 创建委托订单还原引擎

引擎的输入表为标准化后的逐笔委托与逐笔成交合并表,输出表在原始字段基础上增加 orderMark 和 orderIndex 两列,用于标识还原结果来源以及还原后的通道内顺序。

示例脚本

use ops
// 定义输入表:创建逐笔合并数据输入表
colNames = ["TradeDate","SecurityID","Time","SourceType","Type","Price","Qty",
        "BSFlag","BuyNo","SellNo","ApplSeqNum","ChannelNo","Market","ReceiveTime"]
colTypes = ["DATE","SYMBOL","TIME","INT","INT","LONG","LONG","INT",
        "LONG","LONG","LONG","INT","SYMBOL","TIMESTAMP"]
dummyTable = table(1:0, colNames, colTypes)

// 定义输出表:创建委托还原输出表
newColNames = colNames <- ["orderMark", "orderIndex"]
newColTypes = colTypes <-["INT", "LONG"]
try{unsubscribeAll("newOrderTrans")}catch(ex){}
share(streamTable(1:0, newColNames, newColTypes), "newOrderTrans")

// 定义委托还原引擎
inputColMap = 
      { "codeColumn":"SecurityID",    // 证券代码
        "typeColumn":"Type",          // 交易类型
        "priceColumn":"Price",        // 价格
        "qtyColumn":"Qty",            // 数量(股数)
        "buyOrderColumn":"BuyNo",     // 买方订单号
        "sellOrderColumn":"SellNo",   // 卖方订单号
        "sideColumn":"BSFlag",        // 买卖方向
        "msgTypeColumn":"SourceType"  // 数据类型
      }
inputColMap = dict(inputColMap.keys(), inputColMap.values().flatten())
try{dropStreamEngine("OrderReconstitute")}catch(ex){}
engine = createOrderReconstituteEngine(
						name="OrderReconstitute", 
						dummyTable=dummyTable, 
						outputTable=objByName("newOrderTrans"), 
						inputColMap=inputColMap)
  • 定义输入表:

    • colNames 和 colTypes 定义了输入表字段结构,需要与 3.2 节生成的合并表保持一致。

  • 定义输出表:

    • 输出表在输入字段基础上新增 orderMark 和 orderIndex 两列,字段含义见 2.1 小节。

  • 定义委托还原引擎:

    • name:引擎名字。创建引擎前先调用 dropStreamEngine 清理同名引擎,避免重复创建时报错。

    • dummyTable: 指定输入表结构。

    • outputTable:指定引擎输出表。

    • inputColMap:用于指定输入表中各字段在引擎中的语义。详情参考 3.1 节。

创建完成后,后续只需要向引擎追加符合字段结构的上交所合并数据,即可得到还原后的逐笔数据。

4.2 历史批计算调用

进行历史批计算时,可以从 3.2 节保存的上交所合并表 SHOrderTrans 中读取一天一个通道的数据,按 ApplSeqNum 排序后写入委托还原引擎。

示例脚本

// 选取一天一个通道的数据计算
testData = select * from SHOrderTrans
           where TradeDate=2026.04.29 and ChannelNo=1
           order by ApplSeqNum

// 将历史数据追加到委托还原引擎
getStreamEngine("OrderReconstitute").append!(testData)
  • 委托还原引擎会维护每个通道内的订单状态,因此历史批计算时一次只输入同一交易日、同一通道的数据

    • SHOrderTrans 是 3.2 节中提供的示例数据。

    • ApplSeqNum 是交易所在同一个 ChannelNo 内为逐笔委托、逐笔成交等消息分配的连续递增序号,通常从 1 开始递增,用于表示该通道内原始行情消息的发布先后顺序。引擎需要按照这个顺序逐条更新订单簿状态、关联成交与委托号,并识别立即全部成交或立即部分成交等需要还原的场景;如果不同交易日或不同通道的数据混在一起输入,或者同一通道内的 ApplSeqNum 顺序被打乱,就可能导致状态串扰、成交无法正确匹配到对应委托,进而影响还原结果。

    • 后续在快照合成或因子计算中使用上交所数据时,在需要连续顺序列的场景中,要将 ApplSeqNum 字段换成 orderIndex 字段。

  • 调用 getStreamEngine("OrderReconstitute").append!(testData) 后,引擎会将还原后的逐笔数据写入输出表 newOrderTrans。

还原后得到的逐笔数据表如图 4-1 所示。

8. 图 4-1 还原后的逐笔数据表

4.3 实时流计算调用

实时流计算时,可以按照 4.3 节的流程进行前序行情接入和标准化处理,并将逐笔委托、逐笔成交统一写入流表。

  1. 在 DolphinDB 中创建共享流数据表,用于接收处理好的逐笔数据。表结构要求参考 3.1 章节。

  2. 订阅共享流数据表,将合并后的逐笔数据实时写入委托还原引擎中。

  3. 通过行情 SDK 同时订阅逐笔成交和逐笔委托数据,并将同一个通道的逐笔成交和逐笔委托数据写入同一张共享流数据表。在行情 SDK 的回调函数中解析收到的逐笔成交和逐笔委托数据,并完成类型转换、枚举值映射以及字段填充等数据处理工作。数据处理示例参考 3.2 章节。

示例脚本

// 定义输入流表: 首先,创建共享流数据表,用于接收处理好的逐笔数据
colNames = ["TradeDate","SecurityID","Time","SourceType","Type","Price","Qty",
        "BSFlag","BuyNo","SellNo","ApplSeqNum","ChannelNo","Market","ReceiveTime"]
colTypes = ["DATE","SYMBOL","TIME","INT","INT","LONG","LONG","INT",
        "LONG","LONG","LONG","INT","SYMBOL","TIMESTAMP"]
try{unsubscribeAll("OrderTrans")}catch(ex){}
share(streamTable(1:0, colNames, colTypes), "OrderTrans")
  
// 创建订阅: 其次,订阅共享流数据表,将合并后的逐笔数据实时写入委托还原引擎中。
subscribeTable(
      tableName="OrderTrans", actionName="orderTrans_to_reconstitute", 
      offset=-1, batchSize=10000, throttle=0.01, msgAsTable=true,
      handler=getStreamEngine("OrderReconstitute"))
  
// 实时写入数据: 最后将实时数据写入共享流表
testData = select * from SHOrderTrans
           where TradeDate=2026.04.29 and ChannelNo=1
           order by ApplSeqNum
submitJob("replay", "replay", 
          replay{inputTables=testData, 
                        outputTables=objByName("OrderTrans"), 
                        dateColumn="TradeDate", 
                        timeColumn="Time", 
                        replayRate=1000})
  • 定义输入流表:

    • colNames 和 colTypes 定义了输入表字段结构,需要与 3.2 节生成的合并表保持一致。

    • 表类型需要是流表 streamTable,用于后续发布订阅

  • 创建订阅:

    • subscribeTable 会持续订阅 OrderTrans 中新增的数据,并通过 getStreamEngine("OrderReconstitute") 将消息直接追加到委托订单还原引擎。

  • 实时写入数据:

    • SHOrderTrans 是 3.2 节中提供的示例数据。

    • 对示例数据,使用 replay 函数进行控速回放,以此模拟实时数据写入的效果。

    • 使用 submitJob 函数将回放任务作为后台任务执行。可以通过 getRecentJobs 查看后台任务情况。

  • 引擎处理后的结果会持续写入输出表 newOrderTrans。

  • 实时接入时仍需保证同一引擎只处理同一交易日、同一 ChannelNo 的数据,并确保同一通道内按 ApplSeqNum 顺序写入。若多通道并行接入,建议按 ChannelNo 拆分输入流或为每个通道创建独立任务。

5. 应用场景

前文已经介绍了委托订单还原引擎的原理、输入数据准备以及历史批计算和实时流计算的调用方式。本节进一步说明还原后的逐笔数据如何接入下游业务场景,重点展示其在快照合成和因子计算中的使用方法。

5.1 快照合成

委托订单还原引擎补全上交所缺失的原始委托后,输出表可以直接作为 createOrderBookSnapshotEngine 的输入,用于订单簿还原和自定义快照的生成。

与直接使用原始上交所逐笔数据相比,关键差异在于:快照合成引擎需要按照还原后的逻辑顺序处理逐笔数据,因此创建引擎时应将 inputColMap 中的 seqColumn 映射为还原输出表新增的 orderIndex 列,而不是原始 ApplSeqNum 列。

快照合成引擎的详细使用教程参考 基于逐笔数据合成高频 Orderbook:DolphinDB Orderbook 引擎

5.1.1 创建快照合成引擎

下面的示例基于第 4 节得到的上交所还原结果表 newOrderTrans,创建一个上交所快照合成引擎(1 秒频率、10 档深度)。

示例脚本

// 定义输入表: 委托还原引擎的输出表, 详见 4.1 节

// 定义输出表: 创建快照合成引擎输出表
outputColMap, outputTableSch = genOutputColumnsForOBSnapshotEngine(basic=true, time=true, depth=(10, true),tradeDetail=false, orderDetail=false, withdrawDetail=false, prevDetail=false)
colNames = outputTableSch.schema().colDefs.name
colTypes = outputTableSch.schema().colDefs.typeString
share(streamTable(1:0, colNames, colTypes), "newSnapshot")

// 定义快照合成引擎: 
prevCloseDict = dict(STRING, DOUBLE) // 昨收价(key标的, value价格, 可按需替换为真实数据)
inputColMap = 
      { "codeColumn":"SecurityID",    // 证券代码
        "timeColumn":"Time",          // 交易时间
        "typeColumn":"Type",          // 交易类型
        "priceColumn":"Price",        // 价格
        "qtyColumn":"Qty",            // 数量(股数)
        "buyOrderColumn":"BuyNo",     // 买方订单号
        "sellOrderColumn":"SellNo",   // 卖方订单号
        "sideColumn":"BSFlag",        // 买卖方向
        "msgTypeColumn":"SourceType", // 数据类型
        "seqColumn":"orderIndex",     // 逐笔数据序号
        "receiveTime":"ReceiveTime"   // 数据接收时间
      }
inputColMap = dict(inputColMap.keys(), inputColMap.values().flatten())
day = 2026.04.29
try { dropStreamEngine("orderbook") } catch(ex) {}
createOrderBookSnapshotEngine(
        name="orderbook",
        exchange="XSHGSTOCK",
        orderbookDepth=10,
        intervalInMilli=1000,
        date=day,
        startTime=09:30:00.000,
        prevClose=prevCloseDict,
        dummyTable=objByName("newOrderTrans"),
        inputColMap=inputColMap,
        outputTable=objByName("newSnapshot"),
        outputColMap=outputColMap,
        orderBySeq=false
)
  • 定义输入表:

    • 委托还原引擎的输出表, 详见 4.1 节。

  • 定义输出表:

    • genOutputColumnsForOBSnapshotEngine 函数可以获取对应输出表结构,其中:

      • basic=true 表示包含快照基础字段

      • time=true 表示包含快照时间字段

      • depth=(10, true) 表示生成 10 档数据并以 array vector 形式输出

  • 定义快照合成引擎:

    • inputColMap 用于指定还原结果表中各字段在快照合成引擎中的含义。对于上交所还原数据,seqColumn 必须映射到 orderIndex,以保证引擎按照“先委托、后成交”的还原后顺序更新订单簿。

    • dummyTable 指向 newOrderTrans,即委托订单还原引擎的输出表。

    • exchange 指定交易所规则。本例使用 "XSHGSTOCK" 表示上交所股票。

    • intervalInMilli 指定快照输出频率,单位为毫秒。本例设置为 1000,表示每 1 秒输出一次快照。

5.1.2 历史批计算调用

历史批计算时,先按交易日和通道筛选上交所合并数据,写入第 4 节创建的委托订单还原引擎;还原完成后,再将还原结果按 orderIndex 顺序写入快照合成引擎。

示例脚本

// 1. 读取一天一个通道的上交所逐笔合并数据,并按原始通道顺序排序
testData = select * from SHOrderTrans
           where TradeDate=2026.04.29 and ChannelNo=1
           order by ApplSeqNum

// 2. 写入委托订单还原引擎, 按照 4.1 章节创建
getStreamEngine("OrderReconstitute").append!(testData)

// 3. 读取还原后的逐笔数据,并按 orderIndex 排序
reconstitutedData = select * from objByName("newOrderTrans") order by orderIndex

// 4. 写入快照合成引擎, 按照 5.1.1 章节创建
getStreamEngine("orderbook").append!(reconstitutedData)
  • 同一个委托还原引擎和快照合成引擎只能处理同一交易日、同一 ChannelNo 的数据,避免跨交易日状态残留或多通道顺序交错。

  • 还原结果中的 ApplSeqNum 仍保留原始行情消息编号,主要用于溯源;快照合成过程中应使用 orderIndex 作为连续处理顺序。

5.1.3 输出结果与验证

合成的快照如图 5-1 示例:

9. 图 5-1 合成的快照

为了确认委托还原不会改变快照语义,可以将两种方式合成的快照进行对比:

  • orderBookSnapshot1:先经过委托订单还原引擎,再使用还原结果合成的快照。

orderBookSnapshot1 = select * from newSnapshot
  • orderBookSnapshot2:直接使用原始上交所逐笔数据合成的快照。

// 运行 5.1.1 章节重新定义快照合成引擎
// 写入原始上交所数据
testData = select *, 1 as orderMark, ApplSeqNum as orderIndex from SHOrderTrans
           where TradeDate=2026.04.29 and ChannelNo=1
           order by ApplSeqNum
getStreamEngine("orderbook").append!(testData)
// 记录新的计算结果
orderBookSnapshot2 = select * from newSnapshot

如果除顺序编号和更新时间等辅助字段外,两张快照表的盘口价格、盘口数量、成交统计等核心字段完全一致,则说明还原后的逐笔数据可以无缝接入快照合成链路。

使用如下代码进行验证,若返回 true,表示两种链路得到的快照结果一致。

// 除顺序编号和更新时间等辅助字段
dropCols = ["lastseq", "UpdateTime1", "UpdateTime2"]
dropColumns!(orderBookSnapshot1, dropCols)
dropColumns!(orderBookSnapshot2, dropCols)
eqObj(orderBookSnapshot1.values(), orderBookSnapshot2.values())

5.1.4 实时流计算调用

实时流计算时,可以将第 5 节委托订单还原引擎的输出流表作为快照合成引擎的上游数据源,通过 subscribeTable 持续订阅还原后的逐笔数据,并将其追加到 6.1.1 小节创建的快照合成引擎中。

示例脚本

// 运行 4.1 章节创建上交所还原引擎
// 运行 5.1.1 章节重新定义快照合成引擎
  
// 订阅委托订单还原引擎输出表,并将还原后的实时逐笔数据注入快照合成引擎
try{unsubscribeTable(tableName="newOrderTrans", actionName="reconstitute_to_orderbook")}catch(ex){}
subscribeTable(tableName="newOrderTrans", actionName="reconstitute_to_orderbook", 
               offset=-1, batchSize=10000, throttle=0.01, msgAsTable=true,
               handler=getStreamEngine("orderbook"))

// 运行 4.3 章节, 创建 OrderTrans 流表, 并通过回放模拟实时数据写入
  • newOrderTrans 表示委托订单还原引擎的输出流表。

  • subscribeTable 会持续订阅还原后的逐笔数据,并通过 getStreamEngine("orderbook") 将数据实时追加到快照合成引擎。

  • 实时链路中无需再次按 orderIndex 排序,但需要保证上游委托还原引擎接收的数据已经按同一交易日、同一 ChannelNo 内的 ApplSeqNum 顺序写入。

  • 快照合成结果会持续写入 5.1.1 小节中指定的输出表 newSnapshot。

5.2 因子计算

在高频因子计算中,通常需要对委托表与成交表进行复杂的关联操作,且需要针对沪深两市编写不同的数据清洗逻辑。深交所可以直接在原始逐笔委托表上统计,而上交所必须先关联逐笔成交表、反查缺失的委托量,才能得到与深交所口径一致的统计结果。

基于委托还原引擎输出的标准化数据流可以将这一差异消除在数据层:同一个因子只需要维护一套计算逻辑,即可在沪深全市场范围内直接复用。

5.2.1 订单流因子

因子来源:Forecasting high‐frequency excess stock returns via data analytics and machine learning

论文一共给出 39 个高频特征,其中基于委托数据的字段有以下 6 个:

表 5-1 基于委托数据的高频特征

No. 指标名称 中文名
A5 Number of buy orders arriving in the last 60 s 过去 60 秒新增买单委托笔数
A6 Number of sell orders arriving in the last 60 s 过去 60 秒新增卖单委托笔数
A7 Quantity of buy orders arriving in the last 60 s 过去 60 秒新增买单委托量
A8 Quantity of sell orders arriving in the last 60 s 过去 60 秒新增卖单委托量
A26 Volatility of buy order quantity in the last 5 min 过去 5 分钟买单委托量的波动率
A27 Volatility of sell order quantity in the last 5 min 过去 5 分钟卖单委托量的波动率

这六个字段必须依赖委托还原引擎,因为它们统计的是全部到达的委托,无论该委托是否被立即撮合。未经还原的上交所逐笔委托表只保留了未被立即成交、仍在盘口排队的那部分委托,直接在原始数据上统计会产生偏差。还原引擎补全缺失的委托记录后,A5 ~ A8、A26、A27 便可以用与深交所完全一致的口径计算。

基于以上 6 个字段,进一步构造 3 类共 6 个衍生因子:

表 5-2 衍生因子表

因子类型 计算公式 含义
订单笔数因子 NO1 = A5 - A6 净买入委托笔数:过去 60 秒买入委托笔数减去卖出委托笔数
NO3 = (A5-A6)/(A5+A6) 委托笔数不平衡率:净买入委托笔数除以买入、卖出委托笔数之和
订单量因子 QO1 = A7-A8 净买入委托量:过去 60 秒买入委托量减去卖出委托量
QO3 = (A7-A8)/(A7+A8) 委托量不平衡率:净买入委托量除以买入、卖出委托量之和
订单量波动率因子 VolO1 = A26 - A27 委托量标准差差值:过去 5 分钟买方单笔委托量的标准差减去卖方单笔委托量的标准差
VolO3 = (A26-A27)/(A26+A27) 委托量标准差不平衡率:委托量标准差差值除以买方、卖方委托量标准差之和

NO1、QO1、VolO1 是买卖双方原始差值。NO3、QO3、VolO3 是对应的归一化版本,取值范围恒在 [−1, 1] 之间:越接近 1,表示委托笔数 / 委托量 / 委托量波动率越明显地由买方主导;越接近 −1 则相反。

5.2.2 因子计算示例

本小节以 5.2.1 节中的订单流因子为例,介绍如何基于委托订单还原引擎的输出结果,对上交所逐笔数据进行历史批计算。

先按照第 4 节的方法,将上交所逐笔合并数据写入委托订单还原引擎,使用还原后的 newOrderTrans 表作为因子计算输入。

数据过滤:

5.2.1 节中的订单流因子是基于新增委托数据进行计算。而 newOrderTrans 是逐笔合并数据,包含逐笔成交和逐笔委托。因此需要先进行逐笔委托数据过滤。

/** 筛选上交所委托记录 */
def filterSseNewOrder(orderReconTb){
    res = select * from orderReconTb where SourceType=0
    return res
}

因子计算函数

根据因子逻辑,可以将其分为 1 分钟频率的因子和 5 分钟频率的因子。

  • 1 分钟频率的因子计算函数

/** 计算过去 1min 窗口内的指标 */
defg calNewOrderFactor1min(BSFlag, Qty){
    // 买方新增委托笔数
    A5 = sum(iif(BSFlag == 1, 1, 0))
    // 卖方新增委托笔数
    A6 = sum(iif(BSFlag == 2, 1, 0))
    // 净买入委托笔数
    NO1 = A5 - A6
    // 委托笔数不平衡率
    NO3 = (A5 - A6) \ (A5 + A6)

    // 买方新增委托量
    A7 = sum(iif(BSFlag == 1, Qty, 0))
    // 卖方新增委托量
    A8 = sum(iif(BSFlag == 2, Qty, 0))
    // 净买入委托量
    QO1 = A7 - A8
    // 委托量不平衡率
    QO3 = (A7 - A8) \ (A7 + A8)
    return A5, A6, NO1, NO3, A7, A8, QO1, QO3
}
  • 5 分钟频率的因子计算函数

/** 计算过去 5min 窗口内的指标 */
defg calNewOrderFactor5min(BSFlag, Qty){
    // 买方新增委托笔数
    A26 = std(iif(BSFlag == 1, Qty, NULL))
    // 卖方新增委托笔数
    A27 = std(iif(BSFlag == 2, Qty, NULL))
    // 委托量标准差差值
    VolO1 = A26 - A27
    // 委托量标准差不平衡率
    VolO3 = (A26 - A27)/(A26 + A27)
    return A26, A27, VolO1, VolO3
}

因子计算示例(批计算)

使用上述的因子计算函数,可以在 SQL 中直接调用,并用 interval 函数进行 1 分钟和 5 分钟的窗口划分。

/** 计算指标 (批计算)*/
def calNewOrderFactorBatch(newOrderTrans){
    // 过滤出新增委托数据
    orderTb = filterSseNewOrder(newOrderTrans)
    // 计算1分钟指标
    factor1Min = 
        select calNewOrderFactor1min(BSFlag, Qty) as `A5`A6`NO1`NO3`A7`A8`QO1`QO3
        from orderTb
        group by TradeDate, SecurityID, interval(Time, 1m, step=1m, label='right') as Time
    // 计算5分钟指标
    factor5Min = 
        select calNewOrderFactor5min(BSFlag, Qty) as `A26`A27`VolO1`VolO3
        from orderTb
        group by TradeDate, SecurityID, interval(Time, 5m, step=1m, label='right') as Time
    // 关联结果表: 因为频率不一样所以会有空值, 空值用前值填充
    factor = lj(factor1Min, factor5Min, `TradeDate`SecurityID`Time)
    return factor
}

// 0. 按照 4.1 章节代码,创建上交所还原引擎
  
// 1. 读取一天一个通道的上交所逐笔合并数据,并按原始通道顺序排序
testData = select * from SHOrderTrans
           where TradeDate=2026.04.29 and ChannelNo=1 and Time >= 09:30:00.000
           order by ApplSeqNum

// 2. 写入委托订单还原引擎
getStreamEngine("OrderReconstitute").append!(testData)

// 3. 读取还原后的逐笔数据,并按 orderIndex 排序
reconstitutedData = select * from objByName("newOrderTrans")

// 4. 订单流因子计算
factorRes = calNewOrderFactorBatch(reconstitutedData)
  • 批计算整体流程:

    • 首先通过 filterSseNewOrder 从 newOrderTrans 中筛选出新增委托记录。由于还原结果表同时包含委托、成交和产品状态等消息,而订单流因子只统计委托到达情况,因此这里仅保留 SourceType=0 的记录。

    • 然后分别计算 1 分钟窗口和 5 分钟窗口指标。calNewOrderFactor1min 统计过去 1 分钟内买卖两侧的新增委托笔数和委托量,并进一步计算 NO1、NO3、QO1、QO3;calNewOrderFactor5min 统计过去 5 分钟内买卖两侧单笔委托量的标准差,并进一步计算 VolO1 和 VolO3。

    • 两个窗口的结果使用 lj 按 TradeDate、SecurityID 和 Time 左连接,最终形成一张宽表结构的因子结果表。

  • 批计算调用逻辑:

    • 示例中先按交易日和通道号筛选 SHOrderTrans,并按 ApplSeqNum 排序后写入委托订单还原引擎。这样可以保证引擎按照交易所原始发布顺序完成委托还原。

    • 还原完成后,从输出表 newOrderTrans 读取还原后的逐笔数据作为因子计算输入。对于上交所数据,因子计算应基于还原后的委托流,而不是直接基于原始逐笔委托表。

    • 最后调用 calNewOrderFactorBatch(reconstitutedData) 计算订单流因子。输出结果中,A5、A6、A7、A8 对应 1 分钟新增委托统计,A26、A27 对应 5 分钟委托量波动统计,其余字段为基于这些基础指标构造的衍生因子。

输出表采用宽表结构,如图 5-2 示例:

10. 图 5-2 批计算因子结果

注:A26、A27、VolO1、VolO3 基于 5 分钟窗口计算,输出结果的前四行窗口数据不足,所以对应值为空。

因子计算示例(流计算)

/** 计算指标 (流计算)*/
def calNewOrderFactorStream(){
    // 定义输出表
    colNames = ["TradeTime", "SecurityID", 
                "A5", "A6", "NO1", "NO3", 
                "A7", "A8", "QO1", "QO3", 
                "A26", "A27", "VolO1", "VolO3"]
    colTypes = ["TIMESTAMP", "SYMBOL", 
                "DOUBLE", "DOUBLE", "DOUBLE", "DOUBLE", 
                "DOUBLE", "DOUBLE", "DOUBLE", "DOUBLE", 
                "DOUBLE", "DOUBLE", "DOUBLE", "DOUBLE"]
    share(streamTable(1:0, colNames, colTypes), "factorStream")
    // 创建时序聚合引擎进行指标计算
    windowSize = [60000, 5 * 60000]       // 窗口长度: 1min 和 5min
    step = 60000                        // 步长
    metrics = [
        [<calNewOrderFactor1min(BSFlag, Qty) as `A5`A6`NO1`NO3`A7`A8`QO1`QO3>], // 1min 指标
        [<calNewOrderFactor5min(BSFlag, Qty) as `A26`A27`VolO1`VolO3>]          // 5min 指标
    ]
    try{dropStreamEngine("aggFactor")}catch(ex){}
    createTimeSeriesEngine(
            name="aggFactor", 
            windowSize=windowSize, step=step, 
            metrics=metrics, 
            dummyTable=objByName("newOrderTrans"), 
            outputTable=objByName("factorStream"), 
            timeColumn=["TradeDate","Time"], 
            keyColumn="SecurityID",
            fill="ffill",
            closed='left')
    // 订阅 newOrderTrans 数据实时写入计算引擎
    filterHandler = def(msg){
        // 过滤出新增委托数据
        orderTb = filterSseNewOrder(msg)
        getStreamEngine("aggFactor").append!(orderTb)
    }
    try{unsubscribeTable(tableName="newOrderTrans", actionName="aggFactor")}catch(ex){}
    subscribeTable(tableName="newOrderTrans", actionName="aggFactor", 
               offset=-1, batchSize=10000, throttle=0.01, msgAsTable=true,
               handler=filterHandler)
}

// 0. 按照 4.1 章节代码,创建上交所还原引擎
// 0. 按照 4.3 章节代码,创建 OrderTrans 流表和订阅
  
// 1.订单流因子计算流部署
calNewOrderFactorStream()

// 2.实时写入数据: 
// 运行 4.3 章节, 创建 OrderTrans 流表, 并通过回放模拟实时数据写入
  • 创建时序聚合引擎:

    • windowSize=[60000, 5 * 60000] 同时定义 1 分钟和 5 分钟两个滑动窗口,step=60000 表示每 1 分钟输出一次结果。

    • fill="ffill" 用于对低频窗口结果做前值填充。由于 5 分钟窗口与 1 分钟步长同时输出,前值填充可以让 5 分钟因子在每个 1 分钟时间点上都有可用值,便于与 1 分钟因子合并为同一张宽表。

  • 流计算调用逻辑:

    • 部署流计算前,需要先按照 4.1 节创建委托订单还原引擎,并按照 4.3 节创建 OrderTrans 输入流表及其到还原引擎的订阅关系。

    • 调用 calNewOrderFactorStream() 后,系统会创建因子输出流表、时序聚合引擎,并建立从 newOrderTrans 到 aggFactor 的订阅链路。

    • 示例中使用 replay 将历史 SHOrderTrans 数据控速回放到 OrderTrans,用于模拟实时行情接入。实际生产环境中,只需要将行情 SDK 或数据源插件处理后的逐笔合并数据持续写入 OrderTrans 即可。

输出表采用宽表结构,如图 5-3 示例:

11. 图 5-3 流计算因子结果

注:A26、A27、VolO1、VolO3 基于 5 分钟窗口计算,由于流引擎在前四分钟会使用窗口内已有数据进行计算,而不要求窗口满 5 分钟,因此前四行的这四个因子也会计算出非空值。

5.2.3 沪深统一计算

完成上交所委托订单还原后,沪深两市的订单流因子可以收敛到同一套计算框架中。

12. 图 5-4 沪深统一计算框架

本节提供深交所测试数据压缩包:SZ20260429_db.zip

其中包含 2026.04.29 当天 2 支深交所股票的完整逐笔委托数据(SZorder)和逐笔成交数据(SZtrans)。

将数据压缩包上传至 DolphinDB 部署服务器并解压缩:

  • 如果将数据文件夹放置在 DolphinDB 节点主目录下,可以直接运行示例脚本。可在 DolphinDB 客户端执行 getHomeDir() 查询节点主目录。例如 /home/user/DolphinDB/server/SZ20260429_db

  • 上传至服务器其他自定义目录,则需修改示例脚本中的 sz_db 变量后,可以运行示例脚本。

按照 3.1 章节中对 BSFlag、SourceType 等字段的枚举值要求,处理深交所数据。

// 读取指定的磁盘目录下的数据,使用 getHomeDir() 拼接得到绝对路径
sz_db = getHomeDir() + "/SZ20260429_db"

// 加载逐笔委托
SZorder = loadTable(sz_db, "SZorder")

// 处理字段名
orderTypeMap = dict(["1", "2", "U"], [1, 2, 3])
testData = select 
                trade_date as TradeDate, 
                security_id as SecurityID,
                trade_time as Time,
                0 as SourceType,
                orderTypeMap[order_type] as Type,
                long(order_price*10000) as Price,
                order_qty as Qty,
                iif(side=="1", 1, iif(side=="2", 2, NULL)) as BSFlag,
                0 as BuyNo,
                0 as SellNo,
                appl_seq as ApplSeqNum,
                channel_no as ChannelNo,
                exchange as Market,
                now() as ReceiveTime,
                0 as orderMark,
                appl_seq as orderIndex
           from SZorder 
           order by appl_seq

统一计算函数

  • 深交所因子计算示例(批计算)

factorRes = calNewOrderFactorBatch(testData)

沪深数据可以使用与 5.2.2 小节相同的批计算函数 calNewOrderFactorBatch 实现订单流因子的统一计算。

输出表采用宽表结构,如图 5-5 示例:

13. 图 5-5 深交所历史批计算因子结果
  • 深交所因子计算示例(流计算)

use ops
  
// 1.创建接收深交所逐笔委托数据的流表 newOrderTrans
colNames = ["SecurityID", "TradeDate", "Time", "BSFlag", "Qty", "SourceType"]
colTypes = ["SYMBOL", "DATE", "TIME", "INT", "LONG", "INT"]
unsubscribeAll("newOrderTrans")
share(streamTable(1:0, colNames, colTypes), "newOrderTrans")
  
// 2.订单流因子计算流部署
calNewOrderFactorStream()

// 3.实时写入数据:
submitJob("replay", "replay", 
          replay{inputTables=testData, 
                        outputTables=objByName("newOrderTrans"), 
                        dateColumn="TradeDate", 
                        timeColumn="Time", 
                        replayRate=1000})

沪深数据可以使用与 5.2.2 小节相同的流计算函数 calNewOrderFactorStream 进行流任务部署,以此实现订单流因子的统一计算。

实盘上时,只需要将深交所的逐笔委托数据,按照相同的字段映射规则处理后,写入 newOrderTrans 流表即可。

输出表采用宽表结构,如图 5-6 示例:

14. 图 5-6 深交所实时流计算因子结果

6. 结论

上交所逐笔数据在委托记录上的两种“缺失”(立即全部成交不记委托、立即部分成交只记剩余量)是同一套因子逻辑在沪深两市需维护两份代码的根本原因。委托订单还原引擎从数据层解决了此问题,接入还原后的委托流即可像深交所原生委托表一样处理,无需额外关联成交表反推缺失委托量。

本教程围绕此能力覆盖两个典型应用场景:

  • 快照合成:还原后的委托流实现盘口重建所需的“先委托、后成交”顺序,只需将 seqColumn 替换为还原表新增的 orderIndex 列,即可复用快照合成引擎逻辑。

  • 因子计算:以订单流不平衡类因子(NO/QO/VolO)为例,展示沪深两市如何用同一套因子计算函数定义,分别接入沪深各自数据源,产出字段结构完全一致的因子表。

通过委托订单还原补全上交所缺失的委托记录,并统一沪深两市的委托数据口径,可以将交易所差异集中在数据还原层处理。完成还原后,上层的快照合成、因子计算等应用无需再针对不同交易所编写独立的数据处理和计算逻辑,只需基于统一的委托数据进行处理。这样既可以复用已有的快照合成和因子计算逻辑,减少重复开发和维护成本,也为后续扩展更多基于委托数据的量化指标提供统一的数据基础。