DStream::mktDataEngine

首发版本:3.00.6.1

语法

DStream::mktDataEngine(referenceDate, mktDataConfig, [historicalData], [engineConfig])

详情

在 Orca 中创建一个市场数据实时构建引擎。该引擎接收上游输入的原始行情数据,根据指定的市场数据配置构建标准化市场数据对象,并将结果继续传递到下游节点。

该接口适用于 FICC 场景下的实时曲线、曲面、汇率等市场数据构建任务。典型场景是:上游 source 持续注入报价数据,mktDataEngine 将报价转换为标准市场数据,再供 pricingEngine 或其他下游节点消费。

参数

referenceDate DATE 类型标量,表示市场数据的参考日期。

mktDataConfig 字典或由字典组成的元组,表示市场数据构建配置。关于配置格式,请参考下文

historicalData 可选参数,历史市场数据,当构建过程中无法从引擎缓存或实时流中获取所需市场数据时,将从该历史数据中获取。可以是:
  • 字典或向量:数据格式参考 instrumentPricer 函数中的 marketData 参数。

  • 自定义函数:参数是 (kind,date,name)。

engineConfig 可选参数,一个字典,用于设置引擎运行配置。字典包含如下键值对:
  • numThreads:可选,整型标量,表示工作线程数,默认 8。

  • maxQueueDepth:可选,整型标量,表示最大队列深度,默认 10,000,000。

  • useSystemTime:可选,布尔标量,表示是否使用系统时间作为事件时间。默认为 true,使用系统时间作为事件时间。

  • timeColumn:可选,字符串标量,指定时间列(NANOTIMESTAMP)作为事件时间。指定该列后,输入数据中需要包含该列。

  • outputTime:可选,布尔标量,表示是否输出事件时间。默认为 false。

返回值

返回一个 DStream 对象。

例子

本示例通过定义一个流图,将输入的原始汇率报价表(fx_in)接入市场数据引擎。引擎会根据预设的资产配置(FxSpotRate),自动将数据转换为后续定价引擎可直接识别的标准化金融对象,并输出到结果表(mkt_out)中。

if (!existsCatalog("orca")) {
    createCatalog("orca")
}
go
use catalog orca
  
// 定义流图:输入 -> 市场数据引擎 -> 输出
fxConfig = {"name": "USDCNY", "type": "FxSpotRate"}
g = createStreamGraph("simple_mkt_graph")
g.source(`fx_in, `type`name`price, [STRING, STRING, DOUBLE])
 .mktDataEngine(2025.01.01, fxConfig)
 .sink("mkt_out")
g.submit()

// 注入一条原始行情数据
fxQuote = table("FxSpot" as type, "USDCNY" as name, 7.12 as price)
appendOrcaStreamTable("fx_in", fxQuote)

// 查看处理后的标准化结果
select * from useOrcaStreamTable("mkt_out", t -> select * from t)

相关函数:DStream::udfEngineDStream::pricingEnginecreateMktDataEngine