首页 > 开源 > StreamSQL v1.3.0 发布:流-流 JOIN,补齐双流窗口化关联

StreamSQL v1.3.0 发布:流-流 JOIN,补齐双流窗口化关联

OSChina资讯 2026-10-09 11:13 5 阅读 查看原文

 

状态化算子三件套的最后一块。 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

语法

JOIN s2 WITHIN 30 SECONDS ON s1.k=s2.k

ON a.ts BETWEEN b.ts-30s AND b.ts+30s

JOIN s2 WITHIN 30 SECONDS ON ...

JoinWindows.of(Duration)(Java API)

时间语义

事件时间(WITH (TIMESTAMP=...))/ 缺省到达时间

事件时间 + 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