DStream::mktDataEngine
首发版本:3.00.6.1
语法
DStream::mktDataEngine(referenceDate, mktDataConfig, [historicalData],
[engineConfig])
详情
在 Orca 中创建一个市场数据实时构建引擎。该引擎接收上游输入的原始行情数据,根据指定的市场数据配置构建标准化市场数据对象,并将结果继续传递到下游节点。
该接口适用于 FICC 场景下的实时曲线、曲面、汇率等市场数据构建任务。典型场景是:上游 source
持续注入报价数据,mktDataEngine 将报价转换为标准市场数据,再供
pricingEngine 或其他下游节点消费。
参数
referenceDate DATE 类型标量,表示市场数据的参考日期。
mktDataConfig 字典或由字典组成的元组,表示市场数据构建配置。关于配置格式,请参考下文。
-
字典或向量:数据格式参考 instrumentPricer 函数中的 marketData 参数。
-
自定义函数:参数是 (kind,date,name)。
-
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::udfEngine、DStream::pricingEngine、createMktDataEngine
