状态化算子三件套的最后一块。 v1.1 交付 Over 分析函数,v1.2 交付 CEP 模式识别,本版补上 stream-stream JOIN:两条实时流按主键关联、且事件时间相差不超过
WITHIN才算匹配。LEFT 无匹配补 NULL 做缺席检测、3+ 流级联、EmitTo多流喂数。全部新增为 additive,不写 WITHIN 的存量查询一行代码不用改。
一、这一版解决什么问题
此前 StreamSQL 的 JOIN 只有一种形态:流-表富化(实时流挂载静态维表)。但边缘场景大量需求是两条实时流之间的关联:
这类需求没有时间界就会状态无限增长(Flink 普通无界 JOIN 的老问题),主流引擎收敛到的唯一有界形态是时间邻近匹配:|L.ts − R.ts| ≤ WITHIN,过窗状态自动回收。v1.3.0 把它做成了一个 WITHIN 关键字,此前要在业务代码里手写缓冲、自己管清理的活,现在一条 SQL 写完。
二、主流引擎对照
|
|
StreamSQL v1.3.0 |
Flink |
ksqlDB |
Kafka Streams |
|---|---|---|---|---|
|
语法 |
|
|
|
|
|
时间语义 |
事件时间( |
事件时间 + watermark |
stream time |
stream time |
|
无匹配 LEFT |
窗口关闭补 NULL |
水位推过区间补 NULL(append 模式) |
窗口关闭补 NULL |
窗口关闭补 NULL |
|
3+ 流 |
:white_check_mark: 左深级联,每级独立 WITHIN |
:white_check_mark: 任意计划 |
:x: |
:x: |
|
更新流 / retract |
:x:(append-only) |
:white_check_mark: |
:white_check_mark: |
:white_check_mark: |
|
状态回收 |
WITHIN 保留期 + 可选上限双闸 |
state TTL |
窗口关闭 |
保留期 |
三、快速上手
SQL 写法
带 WITHIN 就是双流 JOIN;不带就是流表富化(现状不变)。以下拼写全部接受:WITHIN 30 SECONDS、WITHIN (30 SECONDS)、WITHIN '30s'、WITHIN 100 MS,以及把 WITHIN 写在 ON 之后(两处都写会报错)。
SELECT s.deviceId, s.temperature AS temp, v.vibration AS vib
FROM tempStream AS s
JOIN vibrationStream AS v WITHIN 30 SECONDS
ON s.deviceId = v.deviceId
WHERE s.temperature > 75 AND v.vibration > 30
Go 喂数:EmitTo / Emit
ssql := streamsql.New(
// 可选:每侧缓冲的 key / 行数上限(默认 maxKeys=10000、行数无界)。
// WITHIN 本身就是保留期,正常配置无需关心;高基数场景可显式收紧。
streamsql.WithJoinMaxKeys(20000),
streamsql.WithJoinMaxRows(500000),
)
_ = ssql.Execute(`
SELECT c.cmdId, c.deviceId, r.ackCode
FROM cmdStream AS c
LEFT JOIN ackStream AS r WITHIN 10 SECONDS
ON c.cmdId = r.cmdId
WHERE r.ackCode IS NULL`)
ssql.AddSink(func(rows []map[string]any) {
fmt.Printf("%v\n", rows)
})
// 按流名喂入两侧(流名 = FROM / JOIN 后的名字,大小写敏感)
ssql.EmitTo("cmdStream", map[string]any{"cmdId": "c1", "deviceId": "d1"})
ssql.EmitTo("ackStream", map[string]any{"cmdId": "c1", "ackCode": 0})
// 保留单流习惯:Emit(row) 等价于喂 FROM 侧
// ssql.Emit(map[string]any{"cmdId": "c1", "deviceId": "d1"})
// Stop() 时仍未匹配的 LEFT 行会补发 NULL 再退出,不丢告警
defer ssql.Stop()
注意:EmitTo 只能用于 WITHIN JOIN 查询。普通查询(含下面的两段式下游实例)继续用 Emit;用错会收到明确提示:EmitTo is only available for stream JOIN (WITHIN) queries; use Emit for other queries。EmitSync 对双流查询同样返回 error。IsStreamJoinQuery() 可供上层框架路由查询类型。
输出长什么样
显式投影时,匹配行的左右字段平铺成单键:
{"deviceId":"d1","temperature":61}
两个特例会带上左流别名前缀(形如 o.orderId):LEFT 补 NULL 行和级联输出行。右字段照常平铺,没等到的右字段为 null:
{"amount":99,"o.orderId":"o1","orderId":"o1","payId":null}
SELECT * 时右行整体挂别名为键({左字段…, v: {…}})。建议始终显式投影。
四、场景示例
1. 烟感 × 温感:双源信号互证
单看烟感误报多,单看高温正常:同设备 30 秒内先报警再高温才算火情。
SELECT s.deviceId, t.temperature
FROM smokeStream AS s
JOIN tempStream AS t WITHIN 30 SECONDS
ON s.deviceId = t.deviceId
WHERE s.smoke = true AND t.temperature > 60
匹配是即时的:先到的一侧缓存等待,另一侧到达就配对输出。烟感 d1 与高温 61℃ 在 30 秒内先后到达,输出一行 {"deviceId":"d1","temperature":61};同设备的 50℃ 低温被 WHERE 挡掉,不产出。
2. 订单 × 支付回调:挂单监控
2 分钟没等到支付回调就提醒,用 LEFT JOIN + IS NULL:
SELECT o.orderId, o.amount, p.payId
FROM orderStream AS o
LEFT JOIN payStream AS p WITHIN 120 SECONDS
ON o.orderId = p.orderId
WHERE p.payId IS NULL
已支付的订单正常匹配、被 IS NULL 过滤;没等到的订单在窗口关闭时自动补一行 NULL(payId 为 null),只补一次,不会重复告警。指令下发未回执、告警未确认也是同样写法。注意补 NULL 发生在窗口关闭时(处理时间下约在 WITHIN 时长之后),不是立刻出。
3. 库区 + SKU:复合键配对出入库
入库单和出库单必须库区和 SKU 都相同才算一出一入,防止串库:
SELECT i.zone, i.sku, i.qty AS inQty, o.qty AS outQty
FROM inboundStream AS i
JOIN outboundStream AS o WITHIN 60 SECONDS
ON i.zone = o.zone AND i.sku = o.sku
多键 AND 拼接即可,A 库区的出库单不会配到 B 库区的入库单。AS 别名可用来区分两侧同名字段(inQty / outQty)。
4. 电流 × 电压:两条采样通道按事件时间对齐
电流、电压是独立采样通道,各自带时间戳。按数据自身时间对齐,而不是按到达顺序:
SELECT a.deviceId, a.current, v.voltage
FROM currentStream AS a
JOIN voltageStream AS v WITHIN 5 SECONDS
ON a.deviceId = v.deviceId
WITH (TIMESTAMP = 'ts', TIMEUNIT = 'ms')
两条通道相隔 3 秒的采样点(ts 相差 3000ms)正确配对;网络重发导致的到达乱序不影响归属。车载工况流 × GPS 轨迹流按 VIN 对齐同理。
5. 告警 × 确认人:一对多匹配
一条告警 10 秒内被三个人先后确认,会出三行(窗内可重复匹配,一对多):
SELECT a.alarmId, k.ackBy
FROM alarmStream AS a
JOIN ackStream AS k WITHIN 10 SECONDS
ON a.alarmId = k.alarmId
一条告警 × 三条确认,输出三行 {"alarmId":"a1","ackBy":"zhang/li/wang"}。反之多条告警对一条确认(如批量确认)同样各出一行。
6. 下单 × 支付 × 发货:三流级联跟踪
N 条流 = N−1 个二元 JOIN 级联,每级独立 WITHIN,INNER/LEFT 任意组合:
SELECT o.orderId, p.payId, s.shipNo
FROM orderStream AS o
JOIN payStream AS p WITHIN 5 SECONDS
ON o.orderId = p.orderId
LEFT JOIN shipStream AS s WITHIN 3 SECONDS
ON o.orderId = s.orderId
前两级匹配后,3 秒内等到了运单就带出 shipNo;没等到的在窗口关闭时补 shipNo: null,已支付未发货的订单就这样被筛出来。级联中间行的时间戳取两侧匹配行的较晚者,保证下一级的距离判断有界。
7. 关联后再聚合:两段式
WITHIN JOIN 后直接接窗口/聚合目前不支持,编译期直接报错。等价做法是串两个实例:JOIN 实例的输出喂给下游聚合实例,JOIN 输出保留原始字段(含时间戳):
join := streamsql.New()
_ = join.Execute(`
SELECT s.deviceId, s.ts, s.temperature, v.vibration
FROM tempStream AS s
JOIN vibStream AS v WITHIN 30 SECONDS
ON s.deviceId = v.deviceId`)
agg := streamsql.New()
_ = agg.Execute(`
SELECT deviceId, avg(temperature) AS avgTemp, count(*) AS cnt
FROM joinOut
GROUP BY deviceId, TumblingWindow('1m')`)
// 下游不是 WITHIN JOIN 查询,只能用 Emit 喂数(EmitTo 仅限 WITHIN JOIN 查询)
join.AddSink(func(rows []map[string]any) {
for _, r := range rows {
agg.Emit(r)
}
})
输出效果:两台设备各匹配一对,下游窗口按设备分组输出 [{"deviceId":"d1","avgTemp":70,"cnt":1}, {"deviceId":"d2","avgTemp":80,"cnt":1}](另带 window_id 可做幂等去重)。在 RuleGo 里同构:JOIN 节点(stream_event 输出)接到下游聚合节点即可。
五、详细更新列表
新特性(feat)
- feat: 流-流 JOIN(WITHIN 时间邻近匹配,INNER/LEFT)
- feat: LEFT JOIN 缺席检测,超窗无匹配补 NULL 行
- feat: 3+ 流左深级联 JOIN,每级独立 WITHIN
- feat: 新 API
EmitTo(name, data)多流喂数 - feat: ON 复合键 AND 拼接 +
WITH (TIMESTAMP)事件时间 - feat: 状态有界:WITHIN 必填 + LRU 键/行数双闸 + 指标
- feat: goleak 泄漏检测接入核心包并进 CI
修复(fix)
- fix: SQL
NOT静默失效导致整条 WHERE 丢行 - fix:
WITH多选项(TIMESTAMP/IDLETIMEOUT)互相覆盖 - fix:
WITHIN(括号未闭合被静默接受 - fix: JOIN 过期行可在懒清理缝隙被迟到匹配
链接
- 文档:https://rulego.cc/pages/streamsql-overview/
- RuleGo 生态:https://rulego.cc
- 安装:
go get github.com/rulego/streamsql