Skip to content

Commit 8c99f80

Browse files
feat(query): 补充数据处理函数与异常诊断 (#31)
* docs(json): 制定常用 JSON 函数支持计划 * feat(json): 补充数据库兼容的 JSON 与数据处理函数 * fix(test): 兼容 Java 8 测试编译 * fix(security): 限制高风险函数参数边界 * fix(security): 支持通过 metadata settings 调整函数限制 * refactor(json): 拆分 JSON 函数实现结构 * fix(function): 支持 count 星号参数 * fix(json): 收紧函数边界语义 * fix(function): 保持参数求值顺序 * fix(agg): 支持星号聚合参数 * feat(function): 完善数据处理函数与查询元数据 * refactor: 优化 * refactor: 优化 * fix(sql): 支持中文别名并补齐PR文件 * test(sql): 补充合并与元数据覆盖 * fix(merge): 降低按键合并默认预取 * fix(merge): 避免按键合并窗口等待 * refactor: 优化 * refactor: 优化 * refactor: 优化 * feat: 增强函数支持与异常诊断 - 补充 JSON、时间、字符串等常用函数与 select 列解析场景 - 引入结构化 ReactorQLException 和中英文资源,统一用户可见错误建议 - 增强 ORDER BY、merge_by_key、窗口、聚合、JSON 等异常与安全限制测试 * fix(sql): 修复异常诊断与覆盖率问题 --------- Co-authored-by: zhou-hao <zh.sqy@qq.com>
1 parent d68bf1d commit 8c99f80

57 files changed

Lines changed: 9432 additions & 137 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

docs/plans/2026-07-02-streaming-join-optimization.md

Lines changed: 429 additions & 0 deletions
Large diffs are not rendered by default.

docs/plans/json-functions-support.md

Lines changed: 227 additions & 0 deletions
Large diffs are not rendered by default.
Lines changed: 101 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,101 @@
1+
# ReactorQL 结构化异常诊断开发计划
2+
3+
## 背景
4+
5+
当前 ReactorQL 大量错误直接抛出 `UnsupportedOperationException`,错误信息通常只有“Unsupported expression”或“参数数量错误”。
6+
这类自由文本不利于用户定位 SQL 问题,也不利于大模型根据错误自动修正查询。
7+
8+
## 目标
9+
10+
- 保留现有 `UnsupportedOperationException` 兼容性,不破坏外部调用方和既有测试。
11+
- 为 SQL 解析、编译、函数参数、资源限制、安全限制等错误提供结构化诊断信息。
12+
- 尽可能提供 SQL 表达式、行号、列号、原因、建议和示例。
13+
- 提供国际化扩展点:错误码、默认文案参数,以及默认中英文资源。
14+
- `suggestion``example` 只面向 SQL 使用者提供修复建议或支持用法,不暴露类名、注册表、解析树、执行器、缓存等实现细节。
15+
16+
## 异常模型
17+
18+
新增 `ReactorQLException extends UnsupportedOperationException`,核心字段:
19+
20+
- `i18nCode`:稳定错误码,例如 `error.reactorql.function_argument_count`
21+
- `i18nArgs`:用于平台或调用方国际化渲染的参数。
22+
- `line` / `column`:SQL 起始位置;能从 JSqlParser AST 或 ParseException 取得时填充。
23+
- `expression`:出错 SQL 表达式或原始 SQL。
24+
- `reason`:具体失败原因。
25+
- `suggestion`:推荐修复方向,只描述支持用法或安全边界。
26+
- `example`:可复制的正确写法。
27+
28+
`getMessage()` 输出默认英文诊断文本,`getLocalizedMessage()` 通过 `ResourceBundle` 读取 `i18n/reactorql/messages_*.properties`
29+
JetLinks 平台侧也可以直接读取 `getI18nCode()` / `getI18nArgs()`,再走平台 `LocaleUtils` 或统一错误响应体系。
30+
31+
## 首批错误码
32+
33+
- `error.reactorql.syntax`:SQL 解析失败。
34+
- `error.reactorql.unsupported_expression`:不支持的 select/value 表达式。
35+
- `error.reactorql.unsupported_condition`:不支持的 where/having 条件。
36+
- `error.reactorql.unsupported_group_expression`:不支持的 group by 表达式。
37+
- `error.reactorql.unsupported_from`:不支持的 from 表达式。
38+
- `error.reactorql.unsupported_flat_map`:不支持的列转行表达式。
39+
- `error.reactorql.function_argument_count`:函数参数数量错误。
40+
- `error.reactorql.invalid_argument`:函数参数值、setting 或安全限制错误。
41+
- `error.reactorql.resource_limit`:输入行数、JSON 文本、输出长度、窗口大小等资源限制错误。
42+
43+
## 位置信息策略
44+
45+
- 编译期表达式:利用 JSqlParser 4.6 的 `ASTNodeAccess#getASTNode()``SimpleNode#jjtGetFirstToken()` 读取行列。
46+
- 解析期错误:从 `ParseException.currentToken.next` 读取行列。
47+
- 当前 `SqlParserUtils.quoteNonAsciiAliases(...)` 可能改写 SQL。第一阶段只保证未改写或位置未受影响的 SQL 能返回准确行列;后续如需精准映射,可为 SQL 改写过程补 offset mapping。
48+
49+
## 落地阶段
50+
51+
### 阶段 1:核心模型和高频入口
52+
53+
已覆盖:
54+
55+
- `DefaultReactorQLMetadata(String sql)`:SQL 解析失败包装成 `ReactorQLException`
56+
- `ValueMapFeature#createMapperNow`:不支持的 value/select 表达式。
57+
- `FilterFeature#createPredicateNow`:不支持的 where/having 条件。
58+
- `FunctionMapFeature`:函数参数缺失和数量错误。
59+
- `DefaultReactorQL`:不支持的 select 表达式和 group by 表达式。
60+
- 常见资源/安全参数错误:正则风险、字符串长度、setting、日期单位、`time_bucket` interval。
61+
62+
### 阶段 2:继续替换重点函数
63+
64+
已覆盖:
65+
66+
- JSONPath 相关安全错误:`JsonFunctionSupport``JsonValueSupport`
67+
- `DateFormatFeature``SingleParameterFunctionMapFeature``CoalesceMapFeature`
68+
- `OrderBySupport` 的排序窗口、topN、setting 类型和值域限制。
69+
- `MergeByKeyFeature` 的参数、setting、排序、重复键和单键行数限制。
70+
- 聚合和列转行的参数错误:`MapAggFeature``CollectListAggFeature``CollectRowAggMapFeature``ArrayValueFlatMapFeature`
71+
- FROM 组合函数和子查询集合操作错误:`zip``combine`、子查询 set operation。
72+
- `if``_window` 等常用函数的参数数量和值域错误。
73+
74+
文案边界:
75+
76+
- `reason` 可以描述失败原因,但不要求用户理解内部扩展点。
77+
- `suggestion` 只给出可执行建议,例如“使用 FROM 子句”“增加 LIMIT”“使用简单 JSONPath”“使用 _window('1m')”。
78+
- `example` 给出 SQL 或 setting 写法,不展示内部类名、方法名、缓存、解析树或执行器。
79+
80+
### 阶段 3:运行时上下文增强
81+
82+
- 对运行期函数求值失败增加 row/parameter 上下文,避免只看到底层 `TypeCastException``ArithmeticException`
83+
- 对响应式链路中传播的异常保留原始 cause 和 SQL 表达式。
84+
- 评估是否为 `TypeCastException` 增加结构化诊断字段,或在 ValueMap/Filter 边界统一包装。
85+
86+
## 测试要求
87+
88+
- 断言异常仍然是 `UnsupportedOperationException` 的子类。
89+
- 断言错误码、表达式、原因、建议、示例字段存在。
90+
- 断言语法错误和表达式错误尽量带行列信息。
91+
- 断言 `getLocalizedMessage()` 可以读取英文和中文资源。
92+
- 断言新增错误码在中英文资源中存在。
93+
- 断言 JSON 非安全路径、ORDER BY 资源限制、日期格式、窗口参数、merge_by_key 参数和运行期排序校验等场景使用结构化异常。
94+
- 断言对外建议不包含实现细节关键词。
95+
- 现有错误场景测试保持通过。
96+
97+
当前验证结果:
98+
99+
- `mvn -q test`:260 tests, 0 failures, 0 errors, 0 skipped。
100+
- JaCoCo:instruction 93.57%,branch 81.43%,line 94.52%,不低于基准分支覆盖率。
101+
- `git diff --check`:通过。

docs/plans/streaming-order-by.md

Lines changed: 79 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,79 @@
1+
# 流式 ORDER BY 设计与测试计划
2+
3+
## 背景
4+
5+
ReactorQL 的 `ORDER BY` 运行在 `Flux` 处理链上。精确的 SQL 全局排序必须等上游结束后才能确定第一条输出,在大量数据、长流或无限流中会形成全量物化,带来高延迟、GC 压力和 OOM 风险。
6+
7+
相同约束也出现在主流流式框架中:Reactor `Flux.sort` 会收集全部元素后再排序;Spark Structured Streaming 不支持对输入流直接排序,因为需要跟踪全部已接收数据;Flink SQL 流式 `ORDER BY` 要求主排序字段为递增时间属性,Top-N 和 Window Top-N 则作为受限排序能力提供;Beam 对无界集合的聚合要求窗口或触发器把数据切分成有限集合。
8+
9+
## 目标
10+
11+
1. 保留已有小数据/有界数据 `ORDER BY` 的全局排序语义。
12+
2. 为无 `LIMIT` 的全局排序提供默认最大物化行数,超过后失败而不是持续占用内存。
13+
3.`ORDER BY ... LIMIT [offset,] rowCount` 使用有界 Top-N 策略,只保留 `offset + rowCount` 条候选,语义等价于全局排序后分页。
14+
4. 提供显式窗口排序设置,用于用户接受“只在窗口内有序”的流式场景。
15+
5. 排序表达式的求值保持响应式组合,不再在 comparator 中同步等待异步值。
16+
17+
## 非目标
18+
19+
- 不实现数据源下推排序;数据库、搜索引擎、时序库等源侧优化由数据源实现或上层查询规划负责。
20+
- 不承诺窗口排序满足 SQL 全局 `ORDER BY` 语义。
21+
- 不引入外部状态存储或磁盘 spill;后续如需要可单独设计。
22+
23+
## 方案
24+
25+
### 默认全局排序
26+
27+
- 新增 setting:`orderBy.maxRows`
28+
- 默认值:`10000`;硬上限:`1000000`
29+
-`LIMIT` 且未开启窗口排序时,最多物化 `orderBy.maxRows` 条记录进行排序;超过立即抛出 `UnsupportedOperationException`
30+
31+
### LIMIT Top-N
32+
33+
- 当 SQL 存在 `LIMIT rowCount``LIMIT offset,rowCount`,计算 Top-N 的候选数量为 `offset + rowCount`
34+
- 候选数量不得超过 `orderBy.maxRows``Integer.MAX_VALUE`
35+
- 使用 `PriorityQueue` 保留当前最优 N 条,空间复杂度从全量 `O(total)` 降到 `O(offset + rowCount)`
36+
- 排序完成后再交给既有 `offset``limit` 阶段,保持当前执行链和结果语义。
37+
38+
### 窗口排序
39+
40+
- 新增 setting:`orderBy.windowSize`
41+
- 默认值:`0` 表示关闭;硬上限:`100000`
42+
- 设置为正数时,按固定条数窗口进行局部排序,只保证窗口内有序。
43+
- 这是显式降级模式,适合实时处理“局部有序足够”的场景,不适合作为 SQL 全局排序替代。
44+
45+
## 边界与安全
46+
47+
- setting 可通过 builder 或 SQL hint 进入 metadata,因此保留硬上限,避免查询侧绕过资源保护。
48+
- `orderBy.maxRows <= 0``orderBy.windowSize < 0`、超过硬上限均在构造阶段失败。
49+
- `LIMIT``OFFSET` 为负数时失败。
50+
- `ORDER BY ... LIMIT 0` 不消费排序候选,直接输出空流。
51+
- `NULLS FIRST/LAST` 按 JSqlParser 的 `OrderByElement.NullOrdering` 显式处理;未指定时,ASC 默认 null first,DESC 默认 null last。
52+
53+
## 测试目标
54+
55+
- 既有 `ORDER BY ASC/DESC` 结果保持不变。
56+
-`LIMIT` 超过 `orderBy.maxRows` 时失败,错误消息包含 setting key。
57+
- `ORDER BY ... LIMIT N` 在输入远大于 N 时只要求 `orderBy.maxRows >= N`,结果等价于全局排序后取前 N。
58+
- `ORDER BY ... LIMIT offset,rowCount` 使用 `offset + rowCount` 作为候选上限,结果正确。
59+
- 动态绑定的 `LIMIT ?` 能在执行时识别并走 Top-N。
60+
- Top-N 支持 DESC、`LIMIT 0`
61+
- 固定 `orderBy.windowSize` 只做窗口内排序,验证局部有序结果。
62+
- 多列排序和 `NULLS LAST` 行为正确。
63+
- 非法 setting 构造失败。
64+
65+
## 验证命令
66+
67+
```bash
68+
mvn -q '-Dtest=ReactorQLTest#testOrderBy*' test
69+
mvn -q test
70+
```
71+
72+
## 参考资料
73+
74+
- Reactor `Flux.sort` API 文档:`sort` 会先收集并排序,长流/无限流可能 OOM,建议用 `window` 切分批次。
75+
- Reactor Reference Guide:大量元素可用 grouping/windowing/buffering 做批次化处理。
76+
- Apache Flink SQL `ORDER BY`:流模式下主排序键必须是递增时间属性。
77+
- Apache Flink SQL Top-N / Window Top-N:连续流用 Top-N 或 Window Top-N 约束排序状态,窗口 Top-N 可在窗口结束后清理中间状态。
78+
- Apache Spark Structured Streaming:流式 Dataset 仅在聚合后且 Complete 输出模式下支持排序。
79+
- Apache Beam Programming Guide:无界 PCollection 通过 windowing 拆成有限逻辑窗口,聚合按窗口处理。

pom.xml

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -200,7 +200,7 @@
200200
<plugin>
201201
<groupId>org.jacoco</groupId>
202202
<artifactId>jacoco-maven-plugin</artifactId>
203-
<version>0.8.7</version>
203+
<version>0.8.13</version>
204204
<configuration>
205205
<excludes>
206206
<exclude>**/ExpressionVisitorAdapter*</exclude>
@@ -291,6 +291,12 @@
291291
<version>4.6</version>
292292
</dependency>
293293

294+
<dependency>
295+
<groupId>com.jayway.jsonpath</groupId>
296+
<artifactId>json-path</artifactId>
297+
<version>2.10.0</version>
298+
</dependency>
299+
294300
<dependency>
295301
<groupId>commons-beanutils</groupId>
296302
<artifactId>commons-beanutils</artifactId>
Lines changed: 87 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,87 @@
1+
/*
2+
* Copyright 2025 JetLinks https://www.jetlinks.cn
3+
*
4+
* Licensed under the Apache License, Version 2.0 (the "License");
5+
* you may not use this file except in compliance with the License.
6+
* You may obtain a copy of the License at
7+
*
8+
* http://www.apache.org/licenses/LICENSE-2.0
9+
*
10+
* Unless required by applicable law or agreed to in writing, software
11+
* distributed under the License is distributed on an "AS IS" BASIS,
12+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
* See the License for the specific language governing permissions and
14+
* limitations under the License.
15+
*/
16+
package org.jetlinks.reactor.ql;
17+
18+
import java.util.Objects;
19+
20+
/**
21+
* ReactorQL 查询列结构。
22+
*
23+
* 保存 SQL 解析阶段得到的 select item 别名、原始文本和类型;该对象不可变,可在 SQL AST 释放后继续使用。
24+
*
25+
* @since 1.0.21
26+
*/
27+
public final class Column {
28+
29+
private final String alias;
30+
private final String origin;
31+
private final SQLType type;
32+
33+
public Column(String alias, String origin, SQLType type) {
34+
this.alias = alias;
35+
this.origin = origin;
36+
this.type = type;
37+
}
38+
39+
/**
40+
* @return 查询结果字段别名;没有单一输出字段名的通配列返回 {@code null}
41+
*/
42+
public String getAlias() {
43+
return alias;
44+
}
45+
46+
/**
47+
* @return select item 原始文本,如 {@code count(1) total}
48+
*/
49+
public String getOrigin() {
50+
return origin;
51+
}
52+
53+
/**
54+
* @return 查询列类型
55+
*/
56+
public SQLType getType() {
57+
return type;
58+
}
59+
60+
@Override
61+
public boolean equals(Object o) {
62+
if (this == o) {
63+
return true;
64+
}
65+
if (!(o instanceof Column)) {
66+
return false;
67+
}
68+
Column column = (Column) o;
69+
return Objects.equals(alias, column.alias)
70+
&& Objects.equals(origin, column.origin)
71+
&& type == column.type;
72+
}
73+
74+
@Override
75+
public int hashCode() {
76+
return Objects.hash(alias, origin, type);
77+
}
78+
79+
@Override
80+
public String toString() {
81+
return "Column{" +
82+
"alias='" + alias + '\'' +
83+
", origin='" + origin + '\'' +
84+
", type=" + type +
85+
'}';
86+
}
87+
}

0 commit comments

Comments
 (0)