构建一分钟 K 线聚合与指标计算流图

本文将通过一个完整示例,演示如何使用 Orca 构建一个流计算图,实现对逐笔交易数据的 1 分钟 K 线聚合和常见技术指标(如 EMA、MACD、KDJ)的实时计算。

整个流程共分为五步,以下所有脚本均在计算节点运行。本示例脚本在集群的 cnode1 节点执行,集群总览如下:

1. 图 1-1 集群总览

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 执行。为确保计算资源隔离,系统不会将任务分发至其他计算组的节点。

2. 图 2-1 流图监控

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")