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

发布于

这次 v1.2.0 的升级是 v1.0.0 后一次比较完整的能力提升。新版本补充了 CEP 模式识别、分析函数、资源边界和可观测性相关能力,使得 StreamSQL 不仅能够进行实时聚合,还能处理更复杂的连续事件和状态判断。

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

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

**CEP(复杂事件处理)** 的关注点不在某一条数据的准确与否,而是“一串连续发生的事件,按什么顺序组合起来才有业务意义”。

举例来说:
1. 设备先升温,再振动,再掉线;
2. 车辆 10 秒内连续出现几次急刹;
3. 某个传感器序列满足 A 后跟 B,再跟 C,才算告警。

这些场景对普通窗口聚合功能来说,难以实现。v1.1.0 开始,StreamSQL 引入了 **MATCH_RECOGNIZE** 模式识别,并在后续版本里补齐了几个关键语义。以下是一些新补充的功能:
- **SUBSET**
- **WITHIN** 主动过期清理
- **FINAL** / **RUNNING**
- 贪婪量词最长匹配与 reluctant 区分
- 压测、回归测试与执行路径优化

例如,如果要识别一个风险:同一台设备先连续升温,再出现高振动,就认为进入故障前兆阶段。可以写成如下 SQL:
```sql
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. 分析函数:变化检测和状态判断更好用了

这一版更加完整地支持边缘侧常见的连续判断场景。结合 **OVER**、**PARTITION BY** 和 **WHEN**,很多原本需要自己维护状态的逻辑,现在可以直接用 SQL 表达节省开发成本。

以下是本版本新增的一些分析函数:
- **lag()** / **latest()**:拿上一条或当前最新值,适合用来检测温度突增等;
- **had_changed()** / **changed_col()** / **changed_cols()**:判断字段有无变化,适合设备状态变化检测;
- **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 静默退化** 等问题。本次升级主要补充了以下能力:
- **分组 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 能同时满足“每分钟统计一次均值”和“某设备连续升温后又振动异常则报警”等需求。

### 为什么适合边缘场景

边缘场景面临的现实问题,不在于功能是否能不断增加,而是“**能否在有限资源下,稳定运行常见的实时处理需求**”。StreamSQL 的优势体现在这里:
- **轻量**:采用纯 Go、无外部依赖、秒级启动,适合部署在网关、边缘节点、容器和嵌入式环境;
- **SQL 友好**:熟悉 SQL 的用户可以轻松上手,省去搭建一整套流处理编程模型的时间;
- **嵌入式**:可作为库集成至现有 Go 服务中,无需独立部署;
- **贴近 IoT 场景**:支持设备流数据、维表富化、乱序处理,以及阈值触发与连续状态判断,都是其强项;
- **可观测**:支持 metrics、结构化日志及按实例隔离,适合在业务进程中长期运行。

### 典型应用场景

它特别适合以下场景:
1. **边缘实时聚合**:按设备、产线、站点进行分钟级统计,只将聚合结果回传云端,节省带宽。
2. **设备数据清洗**:在数据接入平台前完成过滤、字段补齐、异常值剔除、单位换算。
3. **维表富化与路由**:按 **deviceId** 关联设备元信息,再按工厂、型号、区域进行分流或统计;
4. **实时告警与变化检测**:及时监测温度突增、状态切换、以及连续异常等问题;
5. **复杂事件识别**:识别“先高温,再高振动,再掉线”这种单条数据无法直接表达的事件序列。

如果你的业务需求是“**数据持续流入,并需要边缘端立刻计算出结果并做出反应**”,StreamSQL 将非常合适。

### 简单的上手示例

下面这个例子展示了如何将 **v1.0.0** 的流表 JOIN 与 **v1.2.0** 的窗口聚合结合,来处理一个典型的边缘分析任务:**设备流数据关联维表后,按工厂和型号做 1 分钟聚合**。
```go
ssql := streamsql.New()
default 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 表达流式过滤、聚合、富化和事件模式识别。两者结合后,可以在边缘侧搭建数据接入、规则处理与实时分析的链路。

---

## 三、详细更新列表
- feat: **cep** 支持 **MATCH_RECOGNIZE** 模式识别,并补齐 **SUBSET**、**FINAL/RUNNING**、**WITHIN** 主动过期清理、贪婪量词最长匹配与 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: 增补事件时间多时间戳聚合、分析分区语义、窗口聚合组合、求值器探针、网关压测基准,以及分析函数、CEP、流表 JOIN、指标测试。
- test: **stream** 覆盖率 57.4% → 77.4%,补强核心路径回归验证。
- fix: **functions** 修复批量聚合 benchmark 入参按 **b.N** 分配导致的 CI 基准 OOM;**test/cep** 压测 drain 改背压,避免 **dataChan** 满导致超时。

---

## 链接
- **文档**:[https://rulego.cc/pages/streamsql-overview/](https://rulego.cc/pages/streamsql-overview/)
- **RuleGo 生态**:[https://rulego.cc/](https://rulego.cc/)
- **安装**:`go get gitee.com/rulego/streamsql` 或者 `go get github.com/rulego/streamsql`

---

原文链接:[点击查看](https://www.oschina.net/news/501980)

评论

暂无评论。

0.054284s