RuleGo RuleGo
🏠首页
  • 快速入门
  • 规则链
  • 标准组件
  • 扩展组件
  • 自定义组件
  • 可视化
  • RuleGo-Server
  • AOP
  • 触发器
  • 高级主题
  • 性能
  • 标准组件
  • 扩展组件
  • 自定义组件
  • 流式计算
  • 组件市场
  • 概述
  • 快速入门
  • 路由
  • DSL
  • API
  • Options
  • 组件
🔥编辑器 (opens new window)
  • 可视化编辑器 (opens new window)
  • RuleGo-Server (opens new window)
  • 🌊StreamSQL
  • 🤖智能体框架
  • 🌬️GFlow 审批工作流 (opens new window)
  • 🦀TPCLAW 智能体平台 (opens new window)
  • ❓问答

    • FAQ
💖支持
👥加入社区
  • Github (opens new window)
  • Gitee (opens new window)
  • GitCode (opens new window)
  • 更新日志 (opens new window)
  • English
  • 简体中文
🏠首页
  • 快速入门
  • 规则链
  • 标准组件
  • 扩展组件
  • 自定义组件
  • 可视化
  • RuleGo-Server
  • AOP
  • 触发器
  • 高级主题
  • 性能
  • 标准组件
  • 扩展组件
  • 自定义组件
  • 流式计算
  • 组件市场
  • 概述
  • 快速入门
  • 路由
  • DSL
  • API
  • Options
  • 组件
🔥编辑器 (opens new window)
  • 可视化编辑器 (opens new window)
  • RuleGo-Server (opens new window)
  • 🌊StreamSQL
  • 🤖智能体框架
  • 🌬️GFlow 审批工作流 (opens new window)
  • 🦀TPCLAW 智能体平台 (opens new window)
  • ❓问答

    • FAQ
💖支持
👥加入社区
  • Github (opens new window)
  • Gitee (opens new window)
  • GitCode (opens new window)
  • 更新日志 (opens new window)
  • English
  • 简体中文

广告采用随机轮播方式显示 ❤️成为赞助商
  • 概述
  • 快速开始
  • 核心概念
  • SQL参考
  • API参考
  • RuleGo集成
  • 加入社区讨论
  • Schema 校验
  • 分析函数
  • 进阶示例
  • 模式识别(CEP)
  • 流-流 JOIN
  • 函数

  • 案例集锦

    • 案例集锦概述
    • 流表 JOIN 元数据增强
    • 会话窗口与设备在线分析
    • 变更数据捕获案例
    • 滑动窗口与持续检测
    • 数据过滤与转换
    • IoT 温度告警与指标聚合
    • 设备故障模式识别(MATCH_RECOGNIZE)
    • 双流关联与缺席告警(流-流 JOIN)
      • 场景 A:双源信号互证(INNER JOIN)
        • 业务目标
        • SQL
        • 输入与输出
        • 行为说明
      • 场景 B:指令下发未回执告警(LEFT JOIN 缺席检测)
        • 业务目标
        • SQL
        • 输入与输出
        • 行为说明
      • 场景 C:车联网工况 × 轨迹(事件时间对齐)
        • 业务目标
        • SQL
        • 输入与输出
        • 行为说明
      • 场景 D:三流级联(a ⋈ b)LEFT ⋈ c
        • 业务目标
        • SQL
        • 输入与输出
        • 行为说明
目录

双流关联与缺席告警(流-流 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)
1
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
1
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
1
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
1
2
3
4

输出(实跑):

{"deviceId":"d1","temp":80.5,"vib":35.2}
{"deviceId":"d2","temp":76,"vib":41}
1
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
1
2
3
4
5

# 输入与输出

cmdStream:

{"cmdId": "c1", "deviceId": "d1"}          // t+0.0s
{"cmdId": "c2", "deviceId": "d2"}          // t+0.3s
1
2

ackStream:

{"cmdId": "c1", "ackCode": 0}              // t+0.6s(c1 的回执准时到达)
1

输出(实跑,约在 t+10~15 秒出现——sweeper 周期为 WITHIN/2,见行为说明):

{"ackCode":null,"cmdId":"c2","deviceId":"d2"}
1

# 行为说明

  • 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')
1
2
3
4
5

# 输入与输出

canStream(ts 为毫秒级 epoch):

{"vin": "V1", "speed": 72, "ts": 1698700000000}    // 事件时间 t0
{"vin": "V2", "speed": 90, "ts": 1698705000000}    // 事件时间 t0+5000s(另一设备 V2,独立时间轴)
1
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 → 匹配
1
2
3

输出(实跑):

{"lon":116.1,"speed":72,"vin":"V1"}
{"lon":117,"speed":90,"vin":"V2"}
1
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
1
2
3
4

# 输入与输出

aStream:

{"k": "k1", "va": 1}                       // t+0.0s
{"k": "k2", "va": 2}                       // t+0.9s
1
2

bStream:

{"k": "k1", "k2": "x", "vb": "B"}          // t+0.3s
{"k": "k2", "k2": "y", "vb": "B2"}         // t+1.2s
1
2

cStream:

{"k2": "x", "vc": "C1"}                    // t+0.6s
{"k2": "z", "vc": "C3"}                    // t+1.5s(没有 k2=y 的配对)
1
2

输出(实跑;第二行约在工况行到达后 3~4.5 秒补发):

{"va":1,"vb":"B","vc":"C1"}
{"va":2,"vb":"B2","vc":null}
1
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 先流入下级,再由下级补发。
在 GitHub 上编辑此页 (opens new window)
上次更新: 2026/09/26, 14:57:53
设备故障模式识别(MATCH_RECOGNIZE)

← 设备故障模式识别(MATCH_RECOGNIZE)

Theme by Vdoing | Copyright © 2023-2026 RuleGo Team | Apache 2.0 License

  • 跟随系统
  • 浅色模式
  • 深色模式
  • 阅读模式