首页 > 开源 > StreamSQL v1.2.0 发布:边缘 SQL 流处理引擎,补齐 CEP 与分析函数能力

StreamSQL v1.2.0 发布:边缘 SQL 流处理引擎,补齐 CEP 与分析函数能力

OSChina资讯 2026-08-18 13:52 2 阅读 查看原文

 

这是 v1.0.0 之后一次比较完整的能力升级。 这一版补上了 CEP 模式识别、分析函数、资源边界和可观测性相关能力,让 StreamSQL 不只是能做实时聚合,也能处理更复杂的连续事件和状态判断。

一、本次升级最值得关注的三件事

1. CEP:边缘事件序列检测终于有了原生 SQL 表达

CEP(Complex Event Processing,复杂事件处理) 看的不是“某一条数据对不对”,而是“一串连续发生的事件,按什么顺序组合起来才有业务意义”。

对很多边缘场景来说,“一条数据是不是异常”并不重要,重要的是“一段连续事件组合起来是不是异常”。例如:

  • 设备先升温,再振动,再掉线;

  • 车辆 10 秒内连续出现几次急刹;

  • 某个传感器序列满足 A 后跟 B,再跟 C,才算真正告警。

这类需求用普通窗口聚合很难写,用业务代码自己维护状态机也很容易变复杂。v1.1.0 开始,StreamSQL 引入了 MATCH_RECOGNIZE 模式识别,并在后续版本里把几个关键语义补齐了:

  • SUBSET

  • WITHIN 主动过期清理

  • FINAL / RUNNING

  • 贪婪量词最长匹配与 reluctant 区分

  • 压测、回归测试与执行路径优化

举个更贴近边缘设备的例子。假设我们要识别这样一种风险:同一台设备先连续升温,再出现高振动,就认为它进入故障前兆阶段。 这时比起“温度超过多少就告警”,更重要的是“这些事件是不是按特定顺序连续发生”。

可以写成这样:

SELECT  FROM stream MATCH_RECOGNIZE (
  PARTITION BY deviceId
  ORDER BY ts
  MEASURES
    A.temperature AS start_temp,
    B.vibration   AS vibration,
    LAST(ts)      AS alarm_time
  PATTERN (A+ B)
  DEFINE
    A AS A.temperature >= 70,
    B AS B.vibration > 30
)

这段 SQL 的意思并不复杂:

  • deviceId 分区,每台设备单独识别;

  • ts 事件时间排序;

  • 如果先出现一段连续高温事件 A+,后面再跟一个高振动事件 B,就输出一次匹配结果;

  • 输出时可以直接取出起始温度、振动值、告警时间等字段。

补完这些之后,StreamSQL 已经可以直接处理单流事件序列识别这类需求,像边缘告警、设备行为检测、规则链前置筛选都会更顺手。

2. 分析函数:变化检测和状态判断更好用了

这一版把边缘侧最常见的连续判断场景补得更完整了。配合 OVERPARTITION BYWHEN,很多原本要自己维护状态的逻辑,现在可以直接写成 SQL。

这一阶段最值得单独拿出来说的函数有:

  • lag() / latest():拿上一条、当前最新值。适合做温度突增、状态跳变、前后差值计算。

  • had_changed() / changed_col() / changed_cols():判断字段有没有变化、哪一列发生了变化。适合做设备状态变化检测、CDC 风格事件提取。

  • hysteresis():做带回滞的阈值判断。适合温度、压力这类会在阈值附近来回抖动的告警场景。

  • latch():做锁存状态。适合“触发后保持告警,直到收到复位信号才解除”这类工业控制场景。

比如:

  • lag(temperature) OVER (PARTITION BY deviceId) 可以判断同一设备相对上一条数据是否突增;

  • had_changed(true, status) 可以筛出状态真正发生变化的时刻,而不是每条状态上报都下游放行;

  • hysteresis(temp, 80, 78) 可以避免温度在 79~80 附近来回抖动时反复报警;

  • latch(set, reset) 可以表达“故障置位后保持,直到复位”为止。

这一轮也补了分析函数在裸 WHERE、复合参数、限定列、分组表达式这些场景下的正确性问题。现在这块能力已经更适合拿来写真实业务规则,而不只是做功能演示。

3. 边缘部署需要的边界,也补得更清楚了

这部分没有 CEP 那么显眼,但对线上更重要。边缘场景更怕的是内存边界不清、窗口追赶异常、写错 SQL 却没有及时暴露

v1.2.0 这一轮,主要把这几件事补清楚了:

  • 分组 LRU 淘汰可观测:补上 group_evicted_count 和节流告警,高基数分组不再悄悄吃内存。

  • 窗口行缓冲上限:新增 WithWindowMaxRows,窗口原始行缓冲终于可以显式设上限。

  • 滑动窗口追赶修复:空闲后恢复流量时,不再长期滞后于数据时间。

  • 解析错误如实上报:子句语法错误不再被吞掉,避免错误 SQL 静默退化成另一条查询。

这些改动的作用也很直接:流量高峰时更容易控内存,排障时更容易看到问题,线上规则写错时也更容易发现。


二、StreamSQL 是什么

一句话定位

StreamSQL 是一个专为物联网边缘场景设计的轻量 SQL 流处理引擎。它用纯 Go 编写、零外部依赖、可嵌入到应用内部运行,用接近 Flink SQL 的写法处理持续不断的数据流。

它要填的空档很明确。传统流处理往往只有两个极端选择:

  • 时序数据库:擅长存储,但实时计算和复杂流式规则表达有限;

  • 重型分布式框架:能力很强,但部署、资源和运维成本对边缘场景来说太重。

StreamSQL 走的是第三条路:把流式过滤、窗口聚合、流表 JOIN、分析函数和 CEP 这些常见能力,放进一个轻量、单机、嵌入式的 Go 引擎里。

核心能力

到 v1.2.0,StreamSQL 已经不只是做时间窗口统计,常见的边缘流处理能力基本都齐了:

  • 实时过滤与转换SELECTWHERE、计算字段、内置函数,适合做清洗、打标、格式转换。

  • 窗口聚合:支持滚动、滑动、会话、计数、全局窗口,覆盖大多数实时统计与阈值触发需求。

  • 流表 JOIN 富化:把设备元信息、工厂、型号、区域等维表字段贴到流上,再继续过滤或聚合。

  • 事件时间处理:支持 watermark、迟到数据、乱序容忍和空闲推进,更适合真实设备数据。

  • 分析函数:支持 OVERPARTITION BY、变化检测、连续分析这类“和上一条、前几条相比”的规则。

  • CEP 模式识别:支持 MATCH_RECOGNIZE,直接描述“先发生 A,再发生 B”这类事件序列。

它既能做“每分钟统计一次均值”,也能做“某设备连续升温后又振动异常则报警”。

为什么适合边缘场景

边缘侧更现实的问题,不是功能能不能一直加,而是“能不能在有限资源下,把常见的实时处理需求稳定跑起来”。

StreamSQL 的优势就在这里:

  • 轻量:纯 Go、零外部依赖、秒级启动,适合网关、边缘节点、容器和嵌入式部署。

  • SQL 友好:会 SQL 的人基本可以直接上手,不需要先搭一整套流处理编程模型。

  • 嵌入式:它是库,不是必须独立部署的大系统,可以直接集成进现有 Go 服务。

  • 场景贴近 IoT:设备流数据、维表富化、乱序时间、阈值触发、连续状态判断,都是它的强项。

  • 可观测:metrics、结构化日志、按实例隔离,适合在业务进程中长期运行。

典型应用场景

它比较适合下面这些场景:

  • 边缘实时聚合:按设备、产线、站点做分钟级统计,只把聚合结果回传云端,节省带宽。

  • 设备数据清洗:在数据接入平台前先完成过滤、字段补齐、异常值剔除、单位换算。

  • 维表富化与路由:按 deviceId 关联设备元信息,再按工厂、型号、区域分流或统计。

  • 实时告警与变化检测:检测温度突增、状态切换、连续异常等问题。

  • 复杂事件识别:识别“先高温,再高振动,再掉线”这类单条数据无法表达的事件序列。

如果你的业务更像“数据持续流入,边缘端要立刻算出结果并做出动作”,那 StreamSQL 会比较对路。

一个上手就能看懂的示例

下面这个例子把 v1.0.0 的流表 JOIN 和 v1.2.0 持续增强的窗口聚合能力放在一起,看一个最典型的边缘分析任务:设备流数据先关联维表,再按工厂和型号做 1 分钟聚合。

ssql := streamsql.New()
defer ssql.Stop()

sql := SELECT m.location, m.model,
               AVG(temperature) AS avg_temp,
               MAX(temperature) AS max_temp,
               COUNT()         AS samples
        FROM stream
        JOIN devices m ON deviceId = m.deviceId
        WHERE temperature > 30
        GROUP BY m.location, m.model, TumblingWindow('1m')

ssql.Execute(sql)

ssql.RegisterTable("devices", []map[string]any{
    {"deviceId": "d1", "location": "plant-A", "model": "X100"},
    {"deviceId": "d2", "location": "plant-B", "model": "X200"},
})

ssql.AddSink(func(results []map[string]any) {
    fmt.Printf("%v\n", results)
})

ssql.Emit(map[string]any{"deviceId": "d1", "temperature": 40.0})
ssql.Emit(map[string]any{"deviceId": "d1", "temperature": 44.0})
ssql.Emit(map[string]any{"deviceId": "d2", "temperature": 35.0})

这个例子里,流里原本只有 deviceId 和温度,但输出结果已经自动带上了工厂和型号。输入是原始设备数据,输出已经是可以直接拿去做业务分析或上报云端的结果。

如果把这个例子再往前推一步:

  • 配合分析函数,可以判断“当前温度相对上一条是否突增”;

  • 配合CEP,可以判断“连续高温之后是否紧接着出现高振动”;

  • 配合事件时间,可以在乱序上报和迟到数据存在时仍得到稳定结果。

它适合谁,不适合谁

适合:物联网边缘计算、设备网关、边缘服务器、本地实时分析与告警、单机或容器内嵌运行、需要在 RuleGo 规则链中增加 SQL 流处理能力的场景。

不适合:需要大规模水平扩展的分布式集群、强依赖持久化状态和事务保障的场景、明显超出单机内存和 CPU 边界的超大吞吐处理。

与 RuleGo 的关系

StreamSQL 是 RuleGo 生态中的流处理基础库。RuleGo 提供输入输出组件和规则编排能力,StreamSQL 负责用 SQL 表达流式过滤、聚合、富化和事件模式识别。两者配合后,可以在边缘侧搭出一条“数据接入 + 规则处理 + SQL 实时分析”的链路。


三、详细更新列表

  • feat: cep 支持 MATCH_RECOGNIZE 模式识别,并补齐 SUBSETFINAL/RUNNINGWITHIN 主动过期 sweeper、贪婪量词最长匹配与 reluctant 区分

  • feat: 分析函数引入 OVER 状态机,支持 PARTITION BY 连接字段、多调用表达式、HAVING 聚合、GROUP BY 表达式

  • feat: functions 新增 hysteresislatch,增强边缘告警抗抖与锁存表达能力

  • feat: aggregator/stream 增加分组 LRU 淘汰可观测能力,新增 group_evicted_count 与节流告警

  • feat: window 增加窗口行缓冲上限 WithWindowMaxRows,丢行可观测

  • fix: JOIN 聚合分组列输出名去表别名前缀,遵循 AS 别名;同名冲突改为编译期报错

  • fix: 修复 row_number/lead 崩溃与静默 nil

  • fix: 未知窗口函数改为明确报错,修复函数表达式参数求值异常,CreateLegacyAggregator panic 改为安全返回

  • fix: watermark 增强防远未来时间戳毒化;datetime 函数支持 time.Time 入参,now() 返回 time.Time,并扩展 to_seconds / unix_timestamp 入参类型

  • fix: JOIN 键做数值归一,GROUP BY 键保留原始类型,减少跨类型匹配异常与分组语义漂移

  • fix: 修复 WHERE 条件中无法调用内置函数,以及分析函数在算术包分析回代、裸 WHERE、复合参数、限定列等场景的运行期静默错值

  • fix: 修复 percentile 聚合第二参数 p 生效问题

  • fix: 修复迟到重复计算、触发竞态、聚合 cast 中断、JOIN panic 崩溃、滑动窗口空闲后长期滞后于数据时间等稳定性问题

  • fix: rsql 子句语法错误如实上报,子句长度不再被静默截断

  • perf: cep 使用 sync.Pool 复用求值基础 map,降低求值分配

  • perf: functions 增加 ListAll 快照缓存,并缓存 usesExprFunction 判定,减少全局锁拷贝与逐行正则开销

  • perf: aggregator/expr/stream 增加分组 LRU 上限、去逐行反射、HAVING 预编译;迟到判重不再依赖 reflect.Pointer

  • refactor: stream 合一同步/异步/窗口前置 per-event pipeline;收敛自定义函数注册框架,补齐单入口测试

  • test: 增补 event-time 多时间戳聚合、分析分区语义、窗口聚合组合、求值器探针、网关压测基准,以及分析函数、CEP、流表 JOIN、指标测试

  • test: stream 覆盖率 57.4% → 77.4%,补强核心路径回归验证

  • fix: functions 修复批量聚合 benchmark 入参按 b.N 分配导致的 CI 基准 OOM;test/cep 压测 drain 改背压,避免 dataChan 满导致超时


链接