这是 v1.0.0 之后一次比较完整的能力升级。 这一版补上了 CEP 模式识别、分析函数、资源边界和可观测性相关能力,让 StreamSQL 不只是能做实时聚合,也能处理更复杂的连续事件和状态判断。
一、本次升级最值得关注的三件事
1. CEP:边缘事件序列检测终于有了原生 SQL 表达
CEP(Complex Event Processing,复杂事件处理) 看的不是“某一条数据对不对”,而是“一串连续发生的事件,按什么顺序组合起来才有业务意义”。
对很多边缘场景来说,“一条数据是不是异常”并不重要,重要的是“一段连续事件组合起来是不是异常”。例如:
这类需求用普通窗口聚合很难写,用业务代码自己维护状态机也很容易变复杂。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 的意思并不复杂:
补完这些之后,StreamSQL 已经可以直接处理单流事件序列识别这类需求,像边缘告警、设备行为检测、规则链前置筛选都会更顺手。
2. 分析函数:变化检测和状态判断更好用了
这一版把边缘侧最常见的连续判断场景补得更完整了。配合 OVER、PARTITION BY 和 WHEN,很多原本要自己维护状态的逻辑,现在可以直接写成 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 已经不只是做时间窗口统计,常见的边缘流处理能力基本都齐了:
-
实时过滤与转换:SELECT、WHERE、计算字段、内置函数,适合做清洗、打标、格式转换。
-
窗口聚合:支持滚动、滑动、会话、计数、全局窗口,覆盖大多数实时统计与阈值触发需求。
-
流表 JOIN 富化:把设备元信息、工厂、型号、区域等维表字段贴到流上,再继续过滤或聚合。
-
事件时间处理:支持 watermark、迟到数据、乱序容忍和空闲推进,更适合真实设备数据。
-
分析函数:支持 OVER、PARTITION 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 模式识别,并补齐 SUBSET、FINAL/RUNNING、WITHIN 主动过期 sweeper、贪婪量词最长匹配与 reluctant 区分
-
feat: 分析函数引入 OVER 状态机,支持 PARTITION BY 连接字段、多调用表达式、HAVING 聚合、GROUP BY 表达式
-
feat: functions 新增 hysteresis 与 latch,增强边缘告警抗抖与锁存表达能力
-
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 满导致超时
链接