构建一分钟 K 线聚合与指标计算流图
本文将通过一个完整示例,演示如何使用 Orca 构建一个流计算图,实现对逐笔交易数据的 1 分钟 K 线聚合和常见技术指标(如 EMA、MACD、KDJ)的实时计算。
整个流程共分为五步,以下所有脚本均在计算节点运行。本示例脚本在集群的 cnode1 节点执行,集群总览如下:
1. 准备工作
增加必要配置项:
- 在配置文件 cluster.cfg、controller.cfg 中配置
enableORCA=true,启动 Orca 功能。 - 在配置文件 cluster.cfg 中配置
persistenceDir,用于持久化公共流表数据;
定义 catalog、指标函数与算子:
if (!existsCatalog("demo")) {
createCatalog("demo")
}
go
use catalog demo
编写技术指标函数:
@state
def EMA(S, N) {
return ::ewmMean(S, span = N, adjust = false)
}
@state
def RD(N, D = 3) {
return ::round(N, D)
}
@state
def HHV(S, N) {
return ::mmax(S, N)
}
@state
def LLV(S, N) {
return ::mmin(S, N)
}
@state
def MACD(CLOSE, SHORT_ = 12, LONG_ = 26, M = 9) {
DIF = EMA(CLOSE, SHORT_) - EMA(CLOSE, LONG_)
DEA = EMA(DIF, M)
MACD = (DIF - DEA) * 2
return RD(DIF, 3), RD(DEA, 3), RD(MACD, 3)
}
@state
def KDJ(CLOSE, HIGH, LOW, N = 9, M1 = 3, M2 = 3) {
RSV = (CLOSE - LLV(LOW, N)) \ (HHV(HIGH, N) - LLV(LOW, N)) * 100
K = EMA(RSV, (M1 * 2 - 1))
D = EMA(K, (M2 * 2 - 1))
J = K * 3 - D * 2
return K, D, J
}
定义聚合算子与指标算子:
aggerators = [
<first(price) as open>,
<max(price) as high>,
<min(price) as low>,
<last(price) as close>,
<sum(volume) as volume>
]
indicators = [
<time>,
<high>,
<low>,
<close>,
<volume>,
<EMA(close, 20) as ema20>,
<EMA(close, 60) as ema60>,
<MACD(close) as `dif`dea`macd>,
<KDJ(close, high, low) as `k`d`j>
]
2. 构建流图并提交
使用 DStream API 定义流图,并提交。
g = createStreamGraph("indicators")
g.source("trade", `time`symbol`price`volume, [DATETIME,SYMBOL,DOUBLE,LONG])
.parallelize(`symbol,3)
.timeSeriesEngine(windowSize=60, step=60, metrics=aggerators, timeColumn=`time, keyColumn=`symbol)
.sync()
.buffer("one_min_bar")
.parallelize(`symbol,2)
.reactiveStateEngine(metrics=indicators, keyColumn=`symbol)
.sync()
.buffer("one_min_indicators")
g.submit()
等待流图状态转为 running 后再进行后续操作,期间可通过以下接口查看流图运行状态:
getStreamGraphMeta()
示例输出(已转为运行中):
| id | fqn | status | semantics | checkpointConfig | ... |
|---|---|---|---|---|---|
| df83...08 | demo.orca_graph.indicators | running | at-least-once | ... | ... |
如图 2-1 所示,在 Web 页面的流图监控模块中可以看到,cnode1 提交的流图在拆分后被分发至同属计算组 group1 的 cnode2 和 cnode3 执行。为确保计算资源隔离,系统不会将任务分发至其他计算组的节点。
3. 插入模拟数据
以下脚本模拟写入了股票的秒级 K 线数据:
sym = symbol(`600519`601398`601288`601857)
ts = 2025.01.01T09:30:00..2025.01.01T15:00:00
n =size(ts) * size(sym)
symbol = stretch(sym, n)
timestamp = take(ts, n)
price = 100+cumsum(rand(0.02, n)-0.01)
volume = rand(1000, n)
trade = table(timestamp, symbol, price, volume)
appendOrcaStreamTable("trade", trade)
4. 查询结果
通过 SQL 获取 1 分钟聚合 K 线和指标:
select * from demo.orca_table.trade
select * from demo.orca_table.one_min_bar
select * from demo.orca_table.one_min_indicators
5. 清理流图
以下脚本会清理流图中定义的所有引擎和流表:
dropStreamGraph("indicators")
