流-流 JOIN
# 流-流 JOIN(Stream-Stream JOIN)
流-流 JOIN 用于把两条实时流按"键相同 + 时间邻近"关联起来(v1.3.0 起支持):ksqlDB 风格的 WITHIN 语法、Flink Interval Join / Kafka Streams JoinWindows 的时间邻近匹配语义、LEFT JOIN 缺席检测,3 条及以上流用左深级联表达。状态有界(WITHIN 即保留期),面向边缘内存定位。
# 与流表 JOIN 的区别
StreamSQL 有两种 JOIN,用是否带 WITHIN 区分:
| 流表 JOIN(Stream-Table) | 流-流 JOIN(Stream-Stream) | |
|---|---|---|
| 关联对象 | 流 × 静态/缓存元数据表 | 流 × 流(两条都是实时数据) |
| 语法标记 | JOIN 表名 ON ...(无 WITHIN) | JOIN 流名 WITHIN 30 SECONDS ON ... |
| 时间约束 | 无(按当前表快照即时挂载) | 必须 \|L.ts − R.ts\| ≤ WITHIN |
| 典型用途 | 补位置/型号等维度属性 | 双源信号互证、指令-回执缺席告警 |
| 注册方式 | RegisterTable | EmitTo 按流名喂入 |
| 文档 | 案例:流表 JOIN 元数据增强 | 本页 + 案例:双流关联与缺席告警 |
漏写 WITHIN 时,JOIN 目标会被当作未注册的表,报 join table %q is not registered 并提示补 WITHIN。
# 快速上手
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
2
3
4
5
ssql := streamsql.New()
_ = ssql.Execute(`...上面的 SQL...`)
ssql.AddSink(func(rows []map[string]interface{}) {
for _, r := range rows {
fmt.Println(r)
}
})
// 按流名喂入两侧数据(流名 = SQL 中 FROM / JOIN 后的名字,大小写敏感)
_ = ssql.EmitTo("tempStream", map[string]interface{}{"deviceId": "d1", "temperature": 80.5})
_ = ssql.EmitTo("vibrationStream", map[string]interface{}{"deviceId": "d1", "vibration": 35.2})
// 停止时,LEFT JOIN 未匹配的 pending 行会先补发 NULL 再退出(Flush)
ssql.Stop()
2
3
4
5
6
7
8
9
10
11
12
13
14
15
习惯单流的代码不用改:Emit(row) 等价于 EmitTo(FROM 流名, row)。
# 在 RuleGo 组件里跑(不用写 Go 代码)
RuleGo 用户无需 Go API——用 x/streamAggregator 组件节点承载 JOIN SQL,streamKey 配置流名来源,多条上游路由汇入同一节点(每条消息的 metadata 带自己的流名):
{"id": "join1", "type": "x/streamAggregator",
"configuration": {
"sql": "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",
"streamKey": "streamName"}}
2
3
4
匹配行经 stream_event 关系流出(Success 是原始数据直通);对关联结果再做窗口聚合用两段式组合,见下文关联后再聚合(两段式)。
# 语法
SELECT s.deviceId, s.temperature, v.vibration
FROM tempStream AS s
JOIN vibrationStream AS v WITHIN 30 SECONDS
ON s.deviceId = v.deviceId
WHERE ...
2
3
4
5
| 语法点 | 规则 |
|---|---|
| WITHIN | 必填,是流-流 JOIN 的语法标记;WITHIN 30 SECONDS / WITHIN (30 SECONDS) / WITHIN '30s' / WITHIN 100 MS 等价接受 |
| WITHIN 位置 | JOIN 表名之后(ksqlDB 位置)或 ON 子句之后,二选一;两处都写报错 |
| JOIN 类型 | INNER JOIN / LEFT JOIN;RIGHT / FULL / CROSS 编译期报错 |
| ON | 等值条件,复合键用 AND 拼接(ON a.x = b.x AND a.y = b.y);非等值报错 |
| 事件时间 | WITH (TIMESTAMP = 'ts') 指定事件时间字段(数值自动判单位 ns/μs/ms/s;小于 10⁹ 的数值按纳秒原样比较——自造序号会永远"邻近",请用真实 epoch 时间戳);不写则按到达时间 |
| 空闲流推进 | WITH (IDLETIMEOUT = '30s'):事件时间模式下,空闲超时的一侧水位推进到墙钟 |
| 输出形状 | 与流表 JOIN 同构:左行字段平铺 + 右行挂别名(v.vibration);别名.字段(左/右皆可)不用 AS 时可能同时输出原名键与去前缀键(f.score → f.score 与 score;级联输出与 LEFT 补 NULL 行同样如此,流表 JOIN 同款行为),建议一律显式 AS;SELECT * 输出 {左字段…, s: {左行}, v: {右行}} |
| 溢出策略 | 沿用 Go API WithOverflowStrategy 配置的 drop/block 策略(SQL WITH 无此选项);expand 不适用于 JOIN 输入(构造期告警并降级 drop) |
WITHIN JOIN 与 CEP 的 WITHIN 是两个东西
MATCH_RECOGNIZE ... WITHIN '1h'(模式识别)约束的是模式序列的活跃期;JOIN ... WITHIN 30 SECONDS(本页)约束的是两行时间戳的距离。二者不共存于一条查询。
# 匹配语义
统一时间模型:每行时间戳 ts = 事件时间(配置了 TIMESTAMP 则必取,缺失按丢行计数)或到达时刻(processing-time)。两种模式走同一条匹配与回收路径。
第一次使用重点看四行
匹配谓词、产出时机、迟到行、LEFT 缺席——这四条决定了你会看到什么输出;其余各行是运维/边界细节,遇到问题再回来查。
| 语义点 | 规则 |
|---|---|
| 匹配谓词 | \|L.ts − R.ts\| ≤ WITHIN(等距区间,同 Kafka Streams JoinWindows / Flink Interval Join) |
| 产出时机 | 到行即匹配即产出(append-only);匹配后不删行,窗内可重复匹配(一对多笛卡尔自然发生) |
| JOIN 键 NULL | 键缺失/NULL 归一为 <nil> 键,且 NULL == NULL 视为相等(沿袭流表 JOIN 约定,与 SQL 标准不同);避免误判请保证 JOIN 键非空 |
| 水位与回收 | 每侧维护 maxSeenTs 最简水位;maxSeenTs − row.ts > WITHIN 的行过期驱逐,已过期的缓冲行不再参与匹配(其出口是驱逐 / LEFT 补 NULL) |
| 迟到行 | 事件时间下,迟到超过 WITHIN 的行被丢弃,计 join_stage{i}_late_dropped 并节流告警(不静默);时钟漂移场景建议 WITHIN 留裕量。MaxOutOfOrderness/AllowedLateness 与 JOIN 组合编译期报错(v1 无乱序裕度配置) |
| 空闲流 | 一侧停发时其水位冻结,LEFT NULL 会延迟;processing-time 模式 sweeper 会把水位推进到墙钟;事件时间模式配 IDLETIMEOUT 后同样推进(权衡:恢复后的迟到行会丢) |
| LEFT 缺席 | 左行到而无匹配 → 挂 pending;窗口关闭(水位推过 / idle 推进 / Stop-Flush)时补一行 NULL(右侧行为空对象挂别名),行级标记保证只补一次 |
| 生命周期 | Stop() 按级联序 Flush:上级 pending LEFT 先流入下级,再逐级补发退出 |
# 3+ 流:左深级联
N 条流 = N−1 个二元 JOIN 按左深级联(SQL 标准结合序,主流引擎同构)。每级独立 WITHIN;INNER/LEFT 任意级组合:
SELECT a.va, b.vb, c.vc
FROM aStream AS a
JOIN bStream AS b WITHIN 5 SECONDS ON a.k = b.k
LEFT JOIN cStream AS c WITHIN 3 SECONDS ON b.k2 = c.k2
2
3
4
- 中间合并行(a⋈b)作为下一级的 LEFT 侧输入,ON 直接引用
a.x/b.y前缀; - 中间行时间戳 = max(两匹配行 ts)(事件时间模式):复合事件的"发生时刻"取最晚组成行,保证下一级距离有界(
|c−b| ≤ W2且|c−a| ≤ W1+W2);processing-time 模式 = 到达本级时刻; - 每级内存闸独立生效,指标带级前缀
join_stage{i}_*(i 从 1 开始)。
# 场景示例
① 双源信号互证(INNER)——同设备 30 秒内"先高温又高振动"才算真异常:
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
2
3
4
5
② 指令下发未回执告警(LEFT 缺席检测)——10 秒内没等到 ACK 的指令(NULL 补齐后由 IS NULL 过滤出来):
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
2
3
4
5
③ 车联网工况 × 轨迹(事件时间)——CAN 流与 GPS 流各自带事件时间戳,按 VIN 关联、5 秒邻近对齐:
SELECT can.vin, can.speed, gps.lon, gps.lat
FROM canStream AS can
JOIN gpsStream AS gps WITHIN 5 SECONDS
ON can.vin = gps.vin
WITH (TIMESTAMP = 'ts')
2
3
4
5
④ 门禁人卡合一(复合键)——刷卡流与抓拍流按 门ID + 人员ID 双键关联,f.score 为人脸比对得分:
SELECT d.gateId, d.userId, f.score
FROM cardStream AS d
JOIN faceStream AS f WITHIN 5 SECONDS
ON d.gateId = f.gateId AND d.userId = f.userId
2
3
4
①②③与三流级联的完整输入/输出实跑记录见案例:双流关联与缺席告警。
# 关联后再聚合(两段式)
WITHIN JOIN 不能与 GROUP BY / 聚合函数写在同一条 SQL(编译期报错)。需要"关联后再统计"(每设备窗口计数、匹配行均值告警等)时,用两段式:JOIN 节点的 stream_event 输出接入下游聚合节点,B 对匹配行做窗口聚合——这也是 RuleGo 组件语境下推荐的组合姿势:
[源A路由] ─┐
├─→ A: x/streamAggregator(JOIN SQL,streamKey=streamName)─stream_event→ B: x/streamAggregator(聚合 SQL)─stream_event→ 告警/存储
[源B路由] ─┘
2
3
规则链 JSON(关键在连线类型):
{
"metadata": {
"nodes": [
{"id": "joinA", "type": "x/streamAggregator", "name": "双流关联",
"configuration": {
"sql": "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",
"streamKey": "streamName"}},
{"id": "aggB", "type": "x/streamAggregator", "name": "窗口聚合",
"configuration": {
"sql": "SELECT deviceId, COUNT(*) AS cnt FROM joinOut GROUP BY deviceId, CountingWindow(10)"}}
],
"connections": [
{"fromId": "joinA", "toId": "aggB", "type": "stream_event"}
]
}
}
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
两点说明(照抄前必读):
streamKey指定消息 metadata 中携带目标流名的键——上例中每条输入消息需带streamName=tempStream或streamName=vibrationStream(键名与值都区分大小写,值必须与 SQL 中 FROM/JOIN 后的流名一致;也支持${}表达式直接解析出流名,如"${msg.stream}");- B 的
FROM joinOut中joinOut是 B 自己的输入流名,可任意命名——B 的数据由 A 经stream_event注入,与这个名字无关。
两条必须遵守的规则(集成测试验证):
- A→B 只能用
stream_event关系连线。A 的Success直通链携带的是原始遥测,连到 B 会把未关联的原始行灌进 B 的窗口(污染计数/均值); - B 用事件时间窗口时,A 的 SQL 必须把时间戳投影出来(如
SELECT ..., s.ts),B 再WITH (TIMESTAMP = 'ts')——JOIN 匹配行不会自动注入时间戳;不投影时 B 按到达时间聚合(CountingWindow或默认窗口)。
B 收到的 stream_event 负载是匹配行数组(每条匹配一条消息),逐行进入 B 自己的流;B 的窗口/聚合语义与单流查询完全一致。单条 SQL 内的"先关联后聚合"仍在路线图(见下方校验清单),两段式是当前的推荐姿势。
# Go API
全部为新增(additive)API,既有 7 个红线 API(New/Execute/Emit/EmitSync/AddSink/Stop/IsAggregationQuery)签名与行为零变化。
ssql := streamsql.New(
streamsql.WithJoinMaxKeys(20000), // 可选:每侧每级缓冲 key 上限(LRU,默认 10000)
streamsql.WithJoinMaxRows(500000), // 可选:每侧每级缓冲行上限(默认 0=无界)
)
err := ssql.Execute(sql) // WITHIN 必填
err = ssql.EmitTo("cmdStream", row) // 按流名喂入;未知名返回 error(含已知流名列表)
ok := ssql.IsStreamJoinQuery() // 判定是否流-流 JOIN 查询(供 RuleGo 组件路由)
2
3
4
5
6
7
Emit(row)=EmitTo(FROM 流名, row);EmitSync对 JOIN 查询返回 error(同 CEP 查询);- 输出经
AddSink/ToChannel获取,与单流查询完全相同。
# 资源边界与内存估算
三道闸逐级兜底:
- WITHIN 必填——保留期即状态上界的时间维;
WithJoinMaxKeys(n)——每侧每级归一键数 LRU 上限(默认 10000)。淘汰含 pending LEFT 行时不补发(计keys_evicted+ 10s 节流告警);过载场景下"淘汰哪一行"不保证确定性(两侧独立入口,调度序不定),正常配置 maxKeys ≥ 峰值活跃键数时零影响;WithJoinMaxRows(n)——每侧每级行缓冲上限(默认 0=无界),超限丢新到行(计rows_dropped+ 节流告警)。
内存估算(每侧每级):
峰值行数 ≈ 输入速率(msg/s) × WITHIN(s) × 活跃 key 占比
峰值内存 ≈ 峰值行数 × ~440 B/行(与窗口缓冲同量级)
级联总内存 ≈ Σ 各级(左+右)峰值;级 i 的左输入速率 = 级 i-1 的输出速率
2
3
实测参考:双流各 1k msg/s、WITHIN=30s → 双侧约 26MB;20 万 distinct key 灌入(maxKeys=10000)→ 缓冲收敛于闸值,全量分配约 109MB 有界。
# 可观测性
GetStats() / Metrics() 暴露级前缀指标(i 为级序,从 1 开始):
| 指标 | 含义 |
|---|---|
join_stage{i}_matches_emitted | 匹配产出行数 |
join_stage{i}_left_timeout_emitted | LEFT 窗口关闭补 NULL 行数 |
join_stage{i}_late_dropped | 迟到 / 缺 TIMESTAMP 字段丢弃行数 |
join_stage{i}_input_dropped | 输入 chan 溢出丢弃(drop 策略或 block 超时) |
join_stage{i}_rows_dropped | maxRows 行缓冲闸丢行 |
join_stage{i}_keys_evicted | maxKeys LRU 淘汰 key 数 |
join_stage{i}_buffered_rows | 当前缓冲行数(gauge) |
join_stage{i}_left_watermark / _right_watermark | 双侧水位(gauge,排查空闲流) |
每类丢弃配 10s 节流 Warn 日志。
# 编译期校验清单
以下组合 Execute 直接报错,无静默忽略:
| 组合 | 处置 |
|---|---|
| WITHIN JOIN + GROUP BY / 窗口 / 聚合函数 | error:单条 SQL 不支持先关联后聚合(用两段式组合,见上文「关联后再聚合」;实际报错原文 WITHIN JOIN cannot be combined with GROUP BY/window/aggregation (...aggregate the JOIN output in a downstream query)) |
| WITHIN JOIN + HAVING / DISTINCT / LIMIT | error:直连尾巴不消费这些子句 |
| WITHIN JOIN + 分析函数(OVER) | error:求值序与双流匹配未定义 |
WITHIN JOIN + MaxOutOfOrderness / AllowedLateness > 0 | error:v1 水位无乱序裕度;IdleTimeout 允许并生效 |
| WITHIN JOIN 与流表 JOIN 混用 | error:v1.3.x 路线("关联后再富化") |
| WITHIN JOIN + unnest | error:行展开仅支持直连路径 |
| RIGHT / FULL / CROSS JOIN | parse error |
| WITHIN 双写(表后 + ON 后) | parse error |
| WITHIN ≤ 0 | parse error |
| FROM/JOIN 流重名 | Execute error |
EmitSync 用于 JOIN 查询 | error |
EmitTo 未知流名 | error(含已知流名列表) |
# 与主流引擎对照
| StreamSQL | Flink | ksqlDB | Kafka Streams | |
|---|---|---|---|---|
| 语法 | JOIN s2 WITHIN 30 SECONDS ON ... | ON a.ts BETWEEN b.ts±30s | JOIN s2 WITHIN 30 SECONDS ON ... | JoinWindows.of(...)(Java API) |
| 时间语义 | 事件时间(TIMESTAMP)/ 缺省到达时间 | 事件时间 + watermark | stream time | stream time |
| 无匹配 LEFT | 窗口关闭补 NULL | 水位推过区间补 NULL | 窗口关闭补 NULL | 窗口关闭补 NULL |
| 更新流 / retract | ❌(全链路 append-only) | ✅ | ✅ | ✅ |
| 3+ 流 | ✅ 左深级联,每级独立 WITHIN | ✅ 任意计划 | ❌ | ❌ |
| 状态回收 | 保留期(WITHIN)+ LRU/行数双闸 | state TTL | 窗口关闭 | 保留期 |
# 路线(v1.3.x / v2)
- RIGHT / FULL JOIN;
ON a.ts BETWEEN b.ts-30s AND b.ts+30s区间谓词语法糖(解析后映射到同一 WITHIN 配置);- 乱序裕度配置(
MaxOutOfOrderness与 JOIN 组合); - 双流 JOIN 输出后再流表富化(WITHIN JOIN 与流表 JOIN 混用);
- RuleGo 组件双流喂数示例(组件仓库发版节奏;拓扑层"多源汇入同一节点"为 RuleGo 既有能力)。