Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
22 commits
Select commit Hold shift + click to select a range
a8ee0aa
docs(json): 制定常用 JSON 函数支持计划
zhou-hao Jul 1, 2026
6f1b3a1
feat(json): 补充数据库兼容的 JSON 与数据处理函数
zhou-hao Jul 1, 2026
95bd0b2
fix(test): 兼容 Java 8 测试编译
zhou-hao Jul 1, 2026
2943735
fix(security): 限制高风险函数参数边界
zhou-hao Jul 1, 2026
bc6f050
fix(security): 支持通过 metadata settings 调整函数限制
zhou-hao Jul 2, 2026
c52143e
refactor(json): 拆分 JSON 函数实现结构
zhou-hao Jul 2, 2026
6a5900d
fix(function): 支持 count 星号参数
zhou-hao Jul 2, 2026
2d6950f
fix(json): 收紧函数边界语义
zhou-hao Jul 2, 2026
0f64a5a
fix(function): 保持参数求值顺序
zhou-hao Jul 2, 2026
51902c8
fix(agg): 支持星号聚合参数
zhou-hao Jul 2, 2026
b4b5260
feat(function): 完善数据处理函数与查询元数据
zhou-hao Jul 2, 2026
9afacd2
refactor: 优化
zhou-hao Jul 2, 2026
d9137aa
refactor: 优化
zhou-hao Jul 2, 2026
47cdb90
fix(sql): 支持中文别名并补齐PR文件
zhou-hao Jul 2, 2026
6df65ae
test(sql): 补充合并与元数据覆盖
zhou-hao Jul 2, 2026
ce4c74c
fix(merge): 降低按键合并默认预取
zhou-hao Jul 2, 2026
3f1af73
fix(merge): 避免按键合并窗口等待
zhou-hao Jul 2, 2026
03703c2
refactor: 优化
zhou-hao Jul 2, 2026
9c55c2e
refactor: 优化
zhou-hao Jul 2, 2026
7cdc1a3
refactor: 优化
zhou-hao Jul 2, 2026
6cb447f
feat: 增强函数支持与异常诊断
zhou-hao Jul 3, 2026
8a70c63
fix(sql): 修复异常诊断与覆盖率问题
zhou-hao Jul 3, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
429 changes: 429 additions & 0 deletions docs/plans/2026-07-02-streaming-join-optimization.md

Large diffs are not rendered by default.

227 changes: 227 additions & 0 deletions docs/plans/json-functions-support.md

Large diffs are not rendered by default.

101 changes: 101 additions & 0 deletions docs/plans/reactorql-exception-diagnostics.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
# ReactorQL 结构化异常诊断开发计划

## 背景

当前 ReactorQL 大量错误直接抛出 `UnsupportedOperationException`,错误信息通常只有“Unsupported expression”或“参数数量错误”。
这类自由文本不利于用户定位 SQL 问题,也不利于大模型根据错误自动修正查询。

## 目标

- 保留现有 `UnsupportedOperationException` 兼容性,不破坏外部调用方和既有测试。
- 为 SQL 解析、编译、函数参数、资源限制、安全限制等错误提供结构化诊断信息。
- 尽可能提供 SQL 表达式、行号、列号、原因、建议和示例。
- 提供国际化扩展点:错误码、默认文案参数,以及默认中英文资源。
- `suggestion` 和 `example` 只面向 SQL 使用者提供修复建议或支持用法,不暴露类名、注册表、解析树、执行器、缓存等实现细节。

## 异常模型

新增 `ReactorQLException extends UnsupportedOperationException`,核心字段:

- `i18nCode`:稳定错误码,例如 `error.reactorql.function_argument_count`。
- `i18nArgs`:用于平台或调用方国际化渲染的参数。
- `line` / `column`:SQL 起始位置;能从 JSqlParser AST 或 ParseException 取得时填充。
- `expression`:出错 SQL 表达式或原始 SQL。
- `reason`:具体失败原因。
- `suggestion`:推荐修复方向,只描述支持用法或安全边界。
- `example`:可复制的正确写法。

`getMessage()` 输出默认英文诊断文本,`getLocalizedMessage()` 通过 `ResourceBundle` 读取 `i18n/reactorql/messages_*.properties`。
JetLinks 平台侧也可以直接读取 `getI18nCode()` / `getI18nArgs()`,再走平台 `LocaleUtils` 或统一错误响应体系。

## 首批错误码

- `error.reactorql.syntax`:SQL 解析失败。
- `error.reactorql.unsupported_expression`:不支持的 select/value 表达式。
- `error.reactorql.unsupported_condition`:不支持的 where/having 条件。
- `error.reactorql.unsupported_group_expression`:不支持的 group by 表达式。
- `error.reactorql.unsupported_from`:不支持的 from 表达式。
- `error.reactorql.unsupported_flat_map`:不支持的列转行表达式。
- `error.reactorql.function_argument_count`:函数参数数量错误。
- `error.reactorql.invalid_argument`:函数参数值、setting 或安全限制错误。
- `error.reactorql.resource_limit`:输入行数、JSON 文本、输出长度、窗口大小等资源限制错误。

## 位置信息策略

- 编译期表达式:利用 JSqlParser 4.6 的 `ASTNodeAccess#getASTNode()` 和 `SimpleNode#jjtGetFirstToken()` 读取行列。
- 解析期错误:从 `ParseException.currentToken.next` 读取行列。
- 当前 `SqlParserUtils.quoteNonAsciiAliases(...)` 可能改写 SQL。第一阶段只保证未改写或位置未受影响的 SQL 能返回准确行列;后续如需精准映射,可为 SQL 改写过程补 offset mapping。

## 落地阶段

### 阶段 1:核心模型和高频入口

已覆盖:

- `DefaultReactorQLMetadata(String sql)`:SQL 解析失败包装成 `ReactorQLException`。
- `ValueMapFeature#createMapperNow`:不支持的 value/select 表达式。
- `FilterFeature#createPredicateNow`:不支持的 where/having 条件。
- `FunctionMapFeature`:函数参数缺失和数量错误。
- `DefaultReactorQL`:不支持的 select 表达式和 group by 表达式。
- 常见资源/安全参数错误:正则风险、字符串长度、setting、日期单位、`time_bucket` interval。

### 阶段 2:继续替换重点函数

已覆盖:

- JSONPath 相关安全错误:`JsonFunctionSupport`、`JsonValueSupport`。
- `DateFormatFeature`、`SingleParameterFunctionMapFeature`、`CoalesceMapFeature`。
- `OrderBySupport` 的排序窗口、topN、setting 类型和值域限制。
- `MergeByKeyFeature` 的参数、setting、排序、重复键和单键行数限制。
- 聚合和列转行的参数错误:`MapAggFeature`、`CollectListAggFeature`、`CollectRowAggMapFeature`、`ArrayValueFlatMapFeature`。
- FROM 组合函数和子查询集合操作错误:`zip`、`combine`、子查询 set operation。
- `if`、`_window` 等常用函数的参数数量和值域错误。

文案边界:

- `reason` 可以描述失败原因,但不要求用户理解内部扩展点。
- `suggestion` 只给出可执行建议,例如“使用 FROM 子句”“增加 LIMIT”“使用简单 JSONPath”“使用 _window('1m')”。
- `example` 给出 SQL 或 setting 写法,不展示内部类名、方法名、缓存、解析树或执行器。

### 阶段 3:运行时上下文增强

- 对运行期函数求值失败增加 row/parameter 上下文,避免只看到底层 `TypeCastException` 或 `ArithmeticException`。
- 对响应式链路中传播的异常保留原始 cause 和 SQL 表达式。
- 评估是否为 `TypeCastException` 增加结构化诊断字段,或在 ValueMap/Filter 边界统一包装。

## 测试要求

- 断言异常仍然是 `UnsupportedOperationException` 的子类。
- 断言错误码、表达式、原因、建议、示例字段存在。
- 断言语法错误和表达式错误尽量带行列信息。
- 断言 `getLocalizedMessage()` 可以读取英文和中文资源。
- 断言新增错误码在中英文资源中存在。
- 断言 JSON 非安全路径、ORDER BY 资源限制、日期格式、窗口参数、merge_by_key 参数和运行期排序校验等场景使用结构化异常。
- 断言对外建议不包含实现细节关键词。
- 现有错误场景测试保持通过。

当前验证结果:

- `mvn -q test`:260 tests, 0 failures, 0 errors, 0 skipped。
- JaCoCo:instruction 93.57%,branch 81.43%,line 94.52%,不低于基准分支覆盖率。
- `git diff --check`:通过。
79 changes: 79 additions & 0 deletions docs/plans/streaming-order-by.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
# 流式 ORDER BY 设计与测试计划

## 背景

ReactorQL 的 `ORDER BY` 运行在 `Flux` 处理链上。精确的 SQL 全局排序必须等上游结束后才能确定第一条输出,在大量数据、长流或无限流中会形成全量物化,带来高延迟、GC 压力和 OOM 风险。

相同约束也出现在主流流式框架中:Reactor `Flux.sort` 会收集全部元素后再排序;Spark Structured Streaming 不支持对输入流直接排序,因为需要跟踪全部已接收数据;Flink SQL 流式 `ORDER BY` 要求主排序字段为递增时间属性,Top-N 和 Window Top-N 则作为受限排序能力提供;Beam 对无界集合的聚合要求窗口或触发器把数据切分成有限集合。

## 目标

1. 保留已有小数据/有界数据 `ORDER BY` 的全局排序语义。
2. 为无 `LIMIT` 的全局排序提供默认最大物化行数,超过后失败而不是持续占用内存。
3. 对 `ORDER BY ... LIMIT [offset,] rowCount` 使用有界 Top-N 策略,只保留 `offset + rowCount` 条候选,语义等价于全局排序后分页。
4. 提供显式窗口排序设置,用于用户接受“只在窗口内有序”的流式场景。
5. 排序表达式的求值保持响应式组合,不再在 comparator 中同步等待异步值。

## 非目标

- 不实现数据源下推排序;数据库、搜索引擎、时序库等源侧优化由数据源实现或上层查询规划负责。
- 不承诺窗口排序满足 SQL 全局 `ORDER BY` 语义。
- 不引入外部状态存储或磁盘 spill;后续如需要可单独设计。

## 方案

### 默认全局排序

- 新增 setting:`orderBy.maxRows`。
- 默认值:`10000`;硬上限:`1000000`。
- 无 `LIMIT` 且未开启窗口排序时,最多物化 `orderBy.maxRows` 条记录进行排序;超过立即抛出 `UnsupportedOperationException`。

### LIMIT Top-N

- 当 SQL 存在 `LIMIT rowCount` 或 `LIMIT offset,rowCount`,计算 Top-N 的候选数量为 `offset + rowCount`。
- 候选数量不得超过 `orderBy.maxRows` 和 `Integer.MAX_VALUE`。
- 使用 `PriorityQueue` 保留当前最优 N 条,空间复杂度从全量 `O(total)` 降到 `O(offset + rowCount)`。
- 排序完成后再交给既有 `offset` 和 `limit` 阶段,保持当前执行链和结果语义。

### 窗口排序

- 新增 setting:`orderBy.windowSize`。
- 默认值:`0` 表示关闭;硬上限:`100000`。
- 设置为正数时,按固定条数窗口进行局部排序,只保证窗口内有序。
- 这是显式降级模式,适合实时处理“局部有序足够”的场景,不适合作为 SQL 全局排序替代。

## 边界与安全

- setting 可通过 builder 或 SQL hint 进入 metadata,因此保留硬上限,避免查询侧绕过资源保护。
- `orderBy.maxRows <= 0`、`orderBy.windowSize < 0`、超过硬上限均在构造阶段失败。
- `LIMIT` 或 `OFFSET` 为负数时失败。
- `ORDER BY ... LIMIT 0` 不消费排序候选,直接输出空流。
- `NULLS FIRST/LAST` 按 JSqlParser 的 `OrderByElement.NullOrdering` 显式处理;未指定时,ASC 默认 null first,DESC 默认 null last。

## 测试目标

- 既有 `ORDER BY ASC/DESC` 结果保持不变。
- 无 `LIMIT` 超过 `orderBy.maxRows` 时失败,错误消息包含 setting key。
- `ORDER BY ... LIMIT N` 在输入远大于 N 时只要求 `orderBy.maxRows >= N`,结果等价于全局排序后取前 N。
- `ORDER BY ... LIMIT offset,rowCount` 使用 `offset + rowCount` 作为候选上限,结果正确。
- 动态绑定的 `LIMIT ?` 能在执行时识别并走 Top-N。
- Top-N 支持 DESC、`LIMIT 0`。
- 固定 `orderBy.windowSize` 只做窗口内排序,验证局部有序结果。
- 多列排序和 `NULLS LAST` 行为正确。
- 非法 setting 构造失败。

## 验证命令

```bash
mvn -q '-Dtest=ReactorQLTest#testOrderBy*' test
mvn -q test
```

## 参考资料

- Reactor `Flux.sort` API 文档:`sort` 会先收集并排序,长流/无限流可能 OOM,建议用 `window` 切分批次。
- Reactor Reference Guide:大量元素可用 grouping/windowing/buffering 做批次化处理。
- Apache Flink SQL `ORDER BY`:流模式下主排序键必须是递增时间属性。
- Apache Flink SQL Top-N / Window Top-N:连续流用 Top-N 或 Window Top-N 约束排序状态,窗口 Top-N 可在窗口结束后清理中间状态。
- Apache Spark Structured Streaming:流式 Dataset 仅在聚合后且 Complete 输出模式下支持排序。
- Apache Beam Programming Guide:无界 PCollection 通过 windowing 拆成有限逻辑窗口,聚合按窗口处理。
8 changes: 7 additions & 1 deletion pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -200,7 +200,7 @@
<plugin>
<groupId>org.jacoco</groupId>
<artifactId>jacoco-maven-plugin</artifactId>
<version>0.8.7</version>
<version>0.8.13</version>
<configuration>
<excludes>
<exclude>**/ExpressionVisitorAdapter*</exclude>
Expand Down Expand Up @@ -291,6 +291,12 @@
<version>4.6</version>
</dependency>

<dependency>
<groupId>com.jayway.jsonpath</groupId>
<artifactId>json-path</artifactId>
<version>2.10.0</version>
</dependency>

<dependency>
<groupId>commons-beanutils</groupId>
<artifactId>commons-beanutils</artifactId>
Expand Down
87 changes: 87 additions & 0 deletions src/main/java/org/jetlinks/reactor/ql/Column.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
/*
* Copyright 2025 JetLinks https://www.jetlinks.cn
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.jetlinks.reactor.ql;

import java.util.Objects;

/**
* ReactorQL 查询列结构。
*
* 保存 SQL 解析阶段得到的 select item 别名、原始文本和类型;该对象不可变,可在 SQL AST 释放后继续使用。
*
* @since 1.0.21
*/
public final class Column {

private final String alias;
private final String origin;
private final SQLType type;

public Column(String alias, String origin, SQLType type) {
this.alias = alias;
this.origin = origin;
this.type = type;
}

/**
* @return 查询结果字段别名;没有单一输出字段名的通配列返回 {@code null}
*/
public String getAlias() {
return alias;
}

/**
* @return select item 原始文本,如 {@code count(1) total}
*/
public String getOrigin() {
return origin;
}

/**
* @return 查询列类型
*/
public SQLType getType() {
return type;
}

@Override
public boolean equals(Object o) {
if (this == o) {
return true;
}
if (!(o instanceof Column)) {
return false;
}
Column column = (Column) o;
return Objects.equals(alias, column.alias)
&& Objects.equals(origin, column.origin)
&& type == column.type;
}

@Override
public int hashCode() {
return Objects.hash(alias, origin, type);
}

@Override
public String toString() {
return "Column{" +
"alias='" + alias + '\'' +
", origin='" + origin + '\'' +
", type=" + type +
'}';
}
}
Loading
Loading