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 的区别
    • 快速上手
      • 在 RuleGo 组件里跑(不用写 Go 代码)
    • 语法
    • 匹配语义
    • 3+ 流:左深级联
    • 场景示例
    • 关联后再聚合(两段式)
    • Go API
    • 资源边界与内存估算
    • 可观测性
    • 编译期校验清单
    • 与主流引擎对照
    • 路线(v1.3.x / v2)
  • 函数

  • 案例集锦

目录

流-流 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
1
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()
1
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"}}
1
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 ...
1
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
1
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
1
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
1
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')
1
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
1
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路由] ─┘
1
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"}
    ]
  }
}
1
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 注入,与这个名字无关。

两条必须遵守的规则(集成测试验证):

  1. A→B 只能用 stream_event 关系连线。A 的 Success 直通链携带的是原始遥测,连到 B 会把未关联的原始行灌进 B 的窗口(污染计数/均值);
  2. 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 组件路由)
1
2
3
4
5
6
7
  • Emit(row) = EmitTo(FROM 流名, row);
  • EmitSync 对 JOIN 查询返回 error(同 CEP 查询);
  • 输出经 AddSink / ToChannel 获取,与单流查询完全相同。

# 资源边界与内存估算

三道闸逐级兜底:

  1. WITHIN 必填——保留期即状态上界的时间维;
  2. WithJoinMaxKeys(n)——每侧每级归一键数 LRU 上限(默认 10000)。淘汰含 pending LEFT 行时不补发(计 keys_evicted + 10s 节流告警);过载场景下"淘汰哪一行"不保证确定性(两侧独立入口,调度序不定),正常配置 maxKeys ≥ 峰值活跃键数时零影响;
  3. WithJoinMaxRows(n)——每侧每级行缓冲上限(默认 0=无界),超限丢新到行(计 rows_dropped + 节流告警)。

内存估算(每侧每级):

峰值行数 ≈ 输入速率(msg/s) × WITHIN(s) × 活跃 key 占比
峰值内存 ≈ 峰值行数 × ~440 B/行(与窗口缓冲同量级)
级联总内存 ≈ Σ 各级(左+右)峰值;级 i 的左输入速率 = 级 i-1 的输出速率
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 既有能力)。
在 GitHub 上编辑此页 (opens new window)
上次更新: 2026/09/26, 14:57:53
模式识别(CEP)
聚合函数

← 模式识别(CEP) 聚合函数→

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

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