# Easy-Agents Federation SQL
基于 Apache Calcite 的 SQL 编译、方言转换、数据源绑定与流式 JDBC 查询底座。
## 模块
- `easy-agents-federation-sql-core`:公共 API、Calcite 编译、单源/联邦自动路由、计划缓存、数据源 Runtime、准入、指标与取消。
- `easy-agents-federation-sql-adapter-jdbc`:默认 JDBC Adapter,也是信创数据库 Adapter 的实现示例。
业务项目通常只需依赖 JDBC Adapter,它会传递依赖 Core:
```xml
com.easyagents
easy-agents-bom
1.2.0-RC
pom
import
com.easyagents
easy-agents-federation-sql-adapter-jdbc
```
## 公共入口
- `FederationSqlEngines.builder()`:组装 Engine、Resolver、策略与准入控制器。
- `engine.sources()`:探测、绑定、预热、更新或移除数据源 Definition。
- `engine.compile()` / `engine.execute()`:高级的节点本地计划模式。
- `engine.query()`:推荐入口,在接收请求的节点完成编译或缓存命中并立即执行。
- `engine.explain()`:显式返回 Calcite 计划;`PHYSICAL` 级别还会请求各数据库的非 `ANALYZE` Explain。
- `engine.cancel(queryId)`:取消准入等待、JDBC 执行或游标消费中的节点本地查询;编译阶段收到取消后不会继续执行。
`FederationSqlPlan` 是 Engine 签发的只读接口,只能交回签发它的 Engine 执行。
`FederationSourceDefinition` 始终描述一个物理数据源。一次查询可见的单源或虚拟联邦范围由调用方使用 `FederationQueryScopeDefinition` 声明;Core 不持久化虚拟数据源,也不保存凭据。
## 最小使用示例
```java
SourceId sourceId = new SourceId("main");
FederationSourceDefinition definition = new FederationSourceDefinition(
sourceId,
1,
JdbcFederationSqlAdapterProvider.ADAPTER_ID,
List.of(new JdbcSchemaDefinition("APP", null, "public")),
Map.of()
);
try (FederationSqlEngine engine = FederationSqlEngines.builder()
.dataSourceResolver(current -> {
HikariDataSource pool = createPool(current.sourceId());
RuntimeFingerprint fingerprint = detectFingerprint(pool);
return FederationDataSourceHandles.owned(pool, fingerprint, pool::close);
})
.maximumPlanCacheEntries(1024)
.maximumPlanCacheWeightBytes(64L * 1024L * 1024L)
.planCacheTimeToLive(Duration.ofMinutes(30))
.build()) {
engine.sources().apply(definition, SourceApplyOptions.prewarmNow());
SqlQueryCommand command = SqlQueryCommand.of(
"SELECT NAME FROM APP.PERSON WHERE ID = ?",
sourceId,
1,
List.of(new SqlParameter(Types.INTEGER, 1))
);
try (FederationResultCursor cursor = engine.query(command)) {
while (cursor.next()) {
System.out.println(cursor.row());
}
}
}
```
`createPool`、凭据存储和 `detectFingerprint` 由调用方实现。Core 管理 Handle/Runtime 生命周期;连接复用、超时、泄漏检测和预热连接数由 HikariCP 等连接池负责。
## 虚拟联邦查询
调用方先分别登记 MySQL 与 PostgreSQL 的物理 `FederationSourceDefinition`,再为一次查询组装逻辑 Binding:
```java
FederationQueryScopeDefinition scope = FederationQueryScopeDefinition.virtual(
"sales-analysis",
7,
Map.of(
"SALES", FederationSourceBindingDefinition.of(
new SourceId("mysql-sales"), 12, Map.of("APP", "APP")
),
"CRM", FederationSourceBindingDefinition.of(
new SourceId("pg-crm"), 5, Map.of("APP", "APP")
)
),
"SALES",
FederationExecutionPolicy.basic()
);
String sql = """
SELECT c.ID, SUM(o.AMOUNT) AS TOTAL
FROM CRM.APP.CUSTOMER c
JOIN SALES.APP.ORDER_ITEM o ON c.ID = o.CUSTOMER_ID
GROUP BY c.ID
ORDER BY TOTAL DESC
""";
try (FederationResultCursor cursor = engine.query(
SqlQueryCommand.of(sql, scope, List.of())
)) {
while (cursor.next()) {
System.out.println(cursor.row());
}
FederationQueryMetricsSnapshot metrics = cursor.metrics();
}
```
查询模式按 Calcite 校验后实际引用的物理 `SourceId` 数量决定。多 Binding Scope 中只引用一个源的 SQL 仍完整下推;引用多个源时,Core 生成目标方言 Fragment,并使用有界 Calcite 本地算子汇总。
调用方已有表列统计快照时,可以通过 `tableStatisticsProvider(...)` 注入行数、行宽、列基数、空值率和唯一键。Provider 的 `snapshot()` 必须一次性返回同时冻结版本、数据和有效期的 `FederationStatisticsSnapshot`,编译阶段不得主动执行 `COUNT(*)`;版本变化会隔离旧计划缓存,计划缓存期限也不会超过统计快照的最早失效时间。统计完整且未过期时,等值 `INNER JOIN` 会把估算搬运量较小的一侧作为本地 Hash Table 构建端;统计缺失、不完整或过期时保持稳定的保守顺序。逻辑 Explain 的每个 Fragment 会返回估算是否可用、扫描/输出行数、行宽、搬运字节、统计来源/采集时间和下推算子。
首批联邦算子覆盖等值 `INNER JOIN`、`LEFT JOIN`、`UNION ALL`、`COUNT/SUM/MIN/MAX/AVG`、普通 `GROUP BY`、CTE、排序和分页。非等值 Join、联邦本地字符比较/排序/分组/`MIN/MAX`、`UNION DISTINCT`、窗口函数、磁盘 Spill 与跨库事务快照会明确拒绝。字符型本地算子需要调用方先统一排序规则,后续再由 Adapter 提供可验证的 Collation 能力。Calcite 本地时间表示只保证毫秒精度;映射精度超过 3 位或运行时检测到亚毫秒值时会明确拒绝。驱动以 `ANY/OTHER` 返回的标准 JDBC 时区标量会保留纳秒并统一为 UTC Offset。
结果采用标准流式 Cursor 语义:Fragment 或本地算子可能在调用方已读取若干行后失败,已交付的行无法撤回。调用方只能在 `next()` 正常返回 `false` 后将本次结果视为完整成功;需要不可逆副作用时应先完整消费并自行提交,或提供补偿机制。
## Explain 与指标
普通 `query` 不会访问数据库 Optimizer。只有显式调用物理 Explain 才会产生额外数据库往返:
```java
SqlCompileRequest compile = SqlCompileRequest.of(sql, scope);
SqlExplainResult logical = engine.explain(
new SqlExplainRequest(compile, SqlExplainLevel.LOGICAL)
);
SqlExplainResult physical = engine.explain(new SqlExplainRequest(compile));
```
`physical.fragments()` 为每个 Fragment 返回目标方言 SQL、参数映射和数据库原生计划。MySQL/PostgreSQL Adapter 尽力归一化扫描方式、候选索引、选中索引、估算行数与过滤条件;数据库没有返回的字段保持空值。为避免原生计划回显敏感常量,Explain 不接受实际参数值,只按 `SqlCompileRequest` 声明的 JDBC 类型绑定 `NULL`,因此索引选择可能与真实参数计划不同。
`FederationResultCursor.metrics()` 可在消费过程中读取,并在耗尽或关闭后定稿,包含模式、计划缓存命中、编译、准入等待、连接获取、数据库执行、本地算子、首行与完整消费耗时,以及最终行/字节、中间搬运行/字节、截断、超时、错误分类和各 Fragment 统计。Adapter 无法安全估算字节时对应字段为 `-1`,不会用 `0` 冒充已测量值。
查询总时限取 Engine、Query Scope 和请求 JDBC timeout 中的最小值。硬时限会覆盖连接池等待后的 JDBC 执行和游标消费,并尝试同时 `cancel`、关闭全部活动 Statement/Cursor;连接池自身仍需配置有限的 connection timeout,以约束 Statement 创建前的连接获取阶段。
## Adapter 扩展
实现 `FederationSqlAdapterProvider` 并通过 Java `ServiceLoader` 注册。Adapter 直接提供 Calcite `Schema`、`SqlDialect`、类型系统、运算符表、Planner Rule 和参数 `SqlDataTypeSpec`,无需额外中间态。重复 `adapterId` 会在启动时拒绝。
## 分布式边界
Definition、revision 和墓碑可以由调用方存入 Redis 等共享状态系统,并通过 `FederationSourceStateProvider` 下发。连接池、Calcite Schema、计划与活动查询均为节点本地对象,不应序列化或跨节点共享。负载均衡请求应携带 `minimumRevision`,落后节点会先同步或返回明确的未就绪错误。
当前联邦路径以最多两个实际物理源和内存内有界汇总为基线。各源使用独立只读连接,不提供跨数据库全局快照一致性;应通过 `FederationExecutionPolicy` 为中间行数、字节数、Fragment 数和总时限设置硬上限。