双流关联与缺席告警(流-流 JOIN)
# 双流关联与缺席告警(流-流 JOIN)
本案例演示 v1.3.0 的流-流 JOIN(WITHIN 语法):两条实时流按"键相同 + 时间邻近"关联,覆盖 INNER 信号互证、LEFT 缺席告警、事件时间对齐、三流级联四个场景。语法与语义详见流-流 JOIN。
运行方式与其他案例不同
流-流 JOIN 有多条输入流,不能用 EmitSync(会报错);FROM 侧可用 Emit,其余流必须按流名 EmitTo 喂入。正确姿势:
ssql := streamsql.New()
_ = ssql.Execute(sql) // 案例的 SQL
ssql.AddSink(func(rows []map[string]interface{}) { /* 收集结果 */ })
_ = ssql.EmitTo("流名", row) // 按流名逐条喂入
ssql.Stop() // 停止时未匹配的 LEFT pending 行补发 NULL(Flush)
2
3
4
5
# 场景 A:双源信号互证(INNER JOIN)
# 业务目标
同一设备有两条上报通道:温度通道和振动通道。单独高温或单独高振动可能是干扰,30 秒内两者同时出现才算真异常。用 INNER JOIN 关联两个通道,双侧 WHERE 各自过滤出异常读数。
# 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
2
3
4
5
# 输入与输出
tempStream:
{"deviceId": "d1", "temperature": 80.5} // t+0.0s
{"deviceId": "d2", "temperature": 76.0} // t+0.3s
{"deviceId": "d3", "temperature": 24.0} // t+0.6s
2
3
vibrationStream:
{"deviceId": "d1", "vibration": 35.2} // t+0.9s
{"deviceId": "d2", "vibration": 41.0} // t+1.2s
{"deviceId": "d9", "vibration": 33.0} // t+1.5s
{"deviceId": "d1", "vibration": 28.0} // t+1.8s
2
3
4
输出(实跑):
{"deviceId":"d1","temp":80.5,"vib":35.2}
{"deviceId":"d2","temp":76,"vib":41}
2
# 行为说明
- d1、d2 两侧条件都满足且在窗内 → 匹配产出;
d3温度 24 被 WHERE 过滤(不产出但已入缓冲);d9只有振动侧、温度侧 30 秒内无同键行 → 不产出; - d1 的第二条振动 28.0 ≤ 30 被 WHERE 过滤;
- 匹配后行不删除:30 秒窗口内 d1 若再来一条合格振动,会和缓冲里的 80.5 再产出一次(窗内可重复匹配);
- 右行挂在别名下:输出形状是左行字段平铺 +
v.vibration引用右行字段。
# 场景 B:指令下发未回执告警(LEFT JOIN 缺席检测)
# 业务目标
平台向设备下发指令,正常情况下设备会在几秒内回 ACK。要抓的是"10 秒内没等到 ACK"的指令——这靠 WHERE 找不出来(不存在的行没有记录),要靠 LEFT JOIN 的"窗口关闭仍无匹配 → 补 NULL 行"语义,再用 IS NULL 把缺席行筛出来。
# SQL
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
# 输入与输出
cmdStream:
{"cmdId": "c1", "deviceId": "d1"} // t+0.0s
{"cmdId": "c2", "deviceId": "d2"} // t+0.3s
2
ackStream:
{"cmdId": "c1", "ackCode": 0} // t+0.6s(c1 的回执准时到达)
输出(实跑,约在 t+10~15 秒出现——sweeper 周期为 WITHIN/2,见行为说明):
{"ackCode":null,"cmdId":"c2","deviceId":"d2"}
# 行为说明
- c1 的 ACK 在窗内到达 → 匹配,且
ackCode=0不满足IS NULL,被 WHERE 过滤——按时回执的指令零输出; - c2 的 ACK 始终没来 → 窗口关闭时补一行
r = {}(空对象),r.ackCode为 NULL,通过 WHERE 产出告警; - processing-time 模式下 sweeper 会把水位推进到墙钟,补发时延 ≈ WITHIN + 一个 sweep 周期(WITHIN/2);
Stop()时未补发的 pending 行也会补发后再退出; - 行级标记保证每条 pending 只补一次,不会重复告警。
# 场景 C:车联网工况 × 轨迹(事件时间对齐)
# 业务目标
CAN 工况流和 GPS 轨迹流各自带设备端打的事件时间戳。两条流网络延迟不同,到达顺序不可靠,必须按事件时间对齐:用 WITH (TIMESTAMP = 'ts') 声明时间字段,按 VIN 关联、事件时间相差 5 秒内才算同一工况片段。
# SQL
SELECT can.vin, can.speed, gps.lon
FROM canStream AS can
JOIN gpsStream AS gps WITHIN 5 SECONDS
ON can.vin = gps.vin
WITH (TIMESTAMP = 'ts')
2
3
4
5
# 输入与输出
canStream(ts 为毫秒级 epoch):
{"vin": "V1", "speed": 72, "ts": 1698700000000} // 事件时间 t0
{"vin": "V2", "speed": 90, "ts": 1698705000000} // 事件时间 t0+5000s(另一设备 V2,独立时间轴)
2
gpsStream:
{"vin": "V1", "lon": 116.1, "ts": 1698700000200} // 距 V1 的 CAN 行 0.2s → 匹配
{"vin": "V1", "lon": 116.2, "ts": 1698700006000} // 距 V1 的 CAN 行 6.0s > WITHIN → 不匹配
{"vin": "V2", "lon": 117.0, "ts": 1698705001000} // 距 V2 的 CAN 行 1.0s → 匹配
2
3
输出(实跑):
{"lon":116.1,"speed":72,"vin":"V1"}
{"lon":117,"speed":90,"vin":"V2"}
2
# 行为说明
- 判定依据是事件时间差,不是到达顺序:第二条 GPS(lon=116.2)虽然紧跟着到达,但事件时间距 CAN 行 6 秒 > WITHIN 5 秒,不匹配——这是事件时间模式与 processing-time 的核心差异;
- 时间戳数值自动判单位(ns/μs/ms/s epoch):本例用毫秒级 epoch(约 1.7×10¹²)被识别为毫秒。注意:小于 10⁹ 的数值视为纳秒原样比较——自造的相对序号(如 1、2、3)之间会永远"邻近",建议业务上使用真实 epoch 时间戳;
- 事件时间模式下若一侧停发,其水位冻结、LEFT 补 NULL 会延迟;配
WITH (IDLETIMEOUT = '30s')可让空闲侧水位推进到墙钟(代价:恢复后的迟到行会被丢); - 迟到超 WITHIN 的行直接丢弃(不参与匹配、不告警,计
join_stage{i}_late_dropped):某侧设备时钟漂移超过 WITHIN 时,该侧行会被系统性丢弃——给 WITHIN 留时钟漂移裕量,或从GetStats()观察该指标。
# 场景 D:三流级联(a ⋈ b)LEFT ⋈ c
# 业务目标
三条流的两级关联:设备状态流 a 与参数流 b 按键 k 关联(5 秒内),合并后的工况再去关联告警流 c(按 k2,3 秒内)。第二级是 LEFT:没有配对告警的工况也要输出(vc 为 NULL),表示"该工况无关联告警"。
# SQL
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
# 输入与输出
aStream:
{"k": "k1", "va": 1} // t+0.0s
{"k": "k2", "va": 2} // t+0.9s
2
bStream:
{"k": "k1", "k2": "x", "vb": "B"} // t+0.3s
{"k": "k2", "k2": "y", "vb": "B2"} // t+1.2s
2
cStream:
{"k2": "x", "vc": "C1"} // t+0.6s
{"k2": "z", "vc": "C3"} // t+1.5s(没有 k2=y 的配对)
2
输出(实跑;第二行约在工况行到达后 3~4.5 秒补发):
{"va":1,"vb":"B","vc":"C1"}
{"va":2,"vb":"B2","vc":null}
2
# 行为说明
- N 条流 = N−1 个二元 JOIN 左深级联:(a⋈b) 的合并行作为第二级的 LEFT 侧输入,ON 直接引用
b.k2前缀; - 第二级 LEFT:
{k2:"z"}的告警行在左侧找不到k2=y的工况(只入缓冲不产出);工况行 (va=2, vb=B2) 窗口关闭无匹配 → 补c = {},c.vc为 NULL; - 中间合并行的时间戳(事件时间模式)= max(两个匹配行 ts)——复合事件的发生时刻取最晚组成行,保证与 c 的距离有界;
- 每级 WITHIN 独立(本级 3 秒不影响上级 5 秒),指标按级前缀区分(
join_stage1_*/join_stage2_*); Stop()按级联序 Flush:上级 pending 先流入下级,再由下级补发。