开头先从实际痛点切入,带出SpringBoot3、Calcite和多数据源这几个核心关键词,然后展开为什么需要这种组合、怎么落地、以及我实测中踩过的坑。
1. 先搞清楚:这个需求到底难在哪
先说个我最近接手的实际场景。业务方要出一张综合报表,数据散落在三个地方:用户主数据在MySQL,订单流水在PostgreSQL,行为日志在ClickHouse。以前的做法是各查各的,然后在Java代码里做内存关联,报表一复杂就写出一堆for循环嵌套。后来单表数据量上来,内存关联慢得没法忍,业务方还时不时提一句“能不能一条SQL把三个库的数据查出来”。
这个诉求听起来简单,实际上是个典型的联邦查询(Federated Query)场景。SpringBoot3里常规的多数据源方案,比如dynamic-datasource,本质上是路由——你指定一个数据源,框架帮你切过去,但一次查询只能落在一个库上,没法在SQL层面做跨库JOIN。ShardingSphere能做,但它的核心优势在分库分表路由和分布式事务,为了一个跨库查询引入一整套中间件,成本和复杂度都不低。
所以我把目光放到了Calcite上。Calcite是Apache旗下的一款SQL解析与查询优化框架,它的定位很特殊:不存储数据、没有自己的存储引擎,但你给它一堆异构数据源,它能帮你做统一SQL解析、校验、优化,甚至把一条SQL下推到各个数据源去执行。SpringBoot3 + Calcite这套组合,相当于在应用层自己做了一个轻量级的“联邦查询引擎”,既能保留多数据源各自的能力,又能用一条SQL把它们串起来。
这篇文章就是我把这套东西从零搭起来、跑通、上线的完整记录,包括依赖选型、Schema构建、SQL执行链路,以及几个让我排查到半夜的坑。下文涉及的代码都是基于我实际运行过的版本整理,你可以直接抄,但建议动手之前先把第一节里几个方案权衡看清楚,不少弯路其实是可以省的。
2. 工程搭建:依赖版本与基础配置
2.1 SpringBoot3项目结构与版本选择
我这次用的是SpringBoot 3.2.4,JDK要求17以上,这没得商量,SpringBoot3强制JDK17起步。如果你还在用JDK8,那第一件事就是先把运行时升上去。项目本身就是一个普通的SpringBoot Web工程,我习惯先建一个干净的骨架,不要急着堆依赖,越干净后面排错越容易。
Calcite这边,我用的版本是1.36.0,在Maven中央仓库直接能拉。要引入的核心依赖是两个:calcite-core和calcite-linq4j。前者提供SQL解析、校验、优化和核心的Schema模型,后者是Calcite内部做表达式计算和数据集遍历用的,很多教程只提core,结果运行时报ClassNotFoundException,这里直接放在pom里。
<dependency> <groupId>org.apache.calcite</groupId> <artifactId>calcite-core</artifactId> <version>1.36.0</version> </dependency> <dependency> <groupId>org.apache.calcite</groupId> <artifactId>calcite-linq4j</artifactId> <version>1.36.0</version> </dependency>顺便提一句,calcite-core会传递依赖一些第三方库,比如guava、jackson、jsr305这些。我在实际项目里遇到过guava版本跟业务代码里其他组件冲突的情况,后面单独开一节说怎么处理,这里先按兵不动。
2.2 多数据源的连接池注册
Calcite本身不负责管理真实的数据源连接,它需要一个能拿到JDBC Connection的入口。我这里统一用HikariCP管理连接池,SpringBoot3默认的DataSource实现也正好是它,所以不需要额外引入别的连接池。
每个数据源在代码里注册成一个独立的HikariDataSource实例,关键是不要让SpringBoot的自动配置把它们当成“唯一数据源”去处理。我的做法是写一个DataSourceRegistry作为注册中心,用一个ConcurrentHashMap保存数据源名称到DataSource实例的映射。
@Component public class DataSourceRegistry { private final Map<String, DataSource> dataSourceMap = new ConcurrentHashMap<>(); public void register(String name, DataSource dataSource) { dataSourceMap.put(name, dataSource); } public DataSource get(String name) { return dataSourceMap.get(name); } public Collection<String> names() { return dataSourceMap.keySet(); } }然后在配置类里把各个数据源初始化好并注册进去。配置项从application.yml里读,包括URL、账号、密码和连接池参数。这里有一个经验:每个数据源务必备上独立的最大连接数和超时时间,因为最慢的那个数据源会拖住整个查询链路,连接池参数不分开调,后面很容易出现某一个源连接耗尽,而其他源空闲的情况。
2.3 依赖冲突的排查经验
Calcite引入后第一道坎往往不是配置问题,而是依赖冲突。我这边的实际经历是guava版本跟项目里另一个人工智能模块的依赖起冲突,启动时直接报NoSuchMethodError。原因很简单:Calcite 1.36.0内部用的是guava 32.x,而业务模块里某个老依赖锁的是guava 27。
这个问题的标准解法是在pom里显式指定guava版本,我用了下面的dependencyManagement把版本钉住:
<dependencyManagement> <dependencies> <dependency> <groupId>com.google.guava</groupId> <artifactId>guava</artifactId> <version>32.1.3-jre</version> </dependency> </dependencies> </dependencyManagement>注意,钉版本之前最好用IDEA的Maven Helper插件看一下依赖树,确认到底是谁引入了低版本guava,避免全局升级之后反过来把别人搞挂。我自己踩过一次:钉住高版本之后,另一个老模块确实启动正常了,但运行到某个序列化方法时开始报错,最后改成那个模块单独排除guava依赖才彻底干净。
3. 核心机制:让Calcite认识你的数据源
这一节是整个集成方案的最关键部分。Calcite之所以能对异构数据源做统一查询,靠的是它自己的元数据模型:Schema、Table和Expression。你首先得把真实数据源的表结构告诉它,它才能在收到SQL的时候知道去哪张表取哪个字段、字段类型是什么。这个“告诉”的过程,就是自定义Schema的实现。
3.1 理解Calcite的三层模型
先解释一下Calcite里的几个核心概念,不然直接看代码容易懵。
Schema是Calcite中最高层的命名空间,一个Schema对应一个库或一个数据源。Table是Schema底下的一张表,描述表名、字段名、字段类型以及如何从底层拿到数据。Calcite的查询流程大致是:收到SQL后先解析,生成AST;然后通过Schema校验表名和字段名;再把逻辑计划发给优化器;优化器借助Table提供的信息生成物理执行计划。整个链路里,Table接口的自定义程度决定了你到底能做多灵活的事情。
我们这次实现的是最常用的方案:用一个AbstractSchema实现类,内部反射式地按表名返回Table,每个Table再映射到真实数据源里的JDBC表。Calcite自身提供了org.apache.calcite.adapter.jdbc.JdbcTable和JdbcSchema,可以直接利用,这也是官方推荐的适配方式。JdbcSchema内部会通过JDBC连接获取DatabaseMetaData来识别表结构,省去了大量手工构建RelDataType的工作。
3.2 自定义Schema实现
我的实现思路是这样:Calcite拿到的Schema名称为multi,下面挂两个子Schema,一个叫orders,对应MySQL中的订单库;一个叫users,对应PostgreSQL中的用户库。每个子Schema背后对应一个真实的数据源。这样做的好处是在SQL里可以直接写multi.orders.t_order,命名空间清晰,扩展新的数据源时只需要注册新Schema,不需要改动查询层代码。
代码上继承Calcite提供的AbstractSchema,在构造时传入数据源名称和数据源连接信息。JdbcSchema有个create方法,传DataSource、Schema名、Dialect即可生成一个可用的Schema实例。我这里直接用它,不需要自己再去解析元数据。
public class MultiSchema extends AbstractSchema { private final String schemaName; private final DataSource dataSource; public MultiSchema(String schemaName, DataSource dataSource) { this.schemaName = schemaName; this.dataSource = dataSource; } @Override protected Map<String, org.apache.calcite.schema.Table> getTableMap() { // 从数据源获取元数据,构建表名到JdbcTable的映射 return buildTableMap(); } }getTableMap的构建逻辑,重点是拿到真实数据库里所有物理表的表名,然后为每个表名创建JdbcTable。JdbcTable的构造需要传入一个JdbcSchema实例和表名。如果数据库表非常多,这一步会比较重,所以我加了一层缓存:应用启动时初始化一次,之后通过刷新接口在表结构变更时手动触发重建。
3.3 在SpringBoot中注册Schema并获取CalciteConnection
Schema构建好之后,需要把它挂到一个Calcite连接上。比较直接的方式是使用Calcite内置的CalciteConnection实现,把RootSchema作为顶层入口,然后在RootSchema下添加子Schema。
我没有用Calcite的Model JSON配置文件方式,而是选择纯Java编程式构建。原因有两点:一是数据源是动态注册的,JSON模型是静态的,程序化构建能对接我们自己DataSourceRegistry;二是出问题时好调试,IDE里直接能看到当前RootSchema挂了几层、每层挂了哪些表。
核心代码大概是这样:
public class CalciteManager { private CalciteConnection calciteConnection; public void init(DataSourceRegistry registry) throws SQLException { Properties info = new Properties(); // Calcite连接需要指定一个内置schema info.setProperty("lex", "JAVA"); Connection connection = DriverManager.getConnection("jdbc:calcite:", info); calciteConnection = connection.unwrap(CalciteConnection.class); RootSchema rootSchema = calciteConnection.getRootSchema(); CalciteSchema calciteSchema = CalciteSchema.createRootSchema(false, false); // 注册orders子Schema SchemaPlus ordersSchema = rootSchema.add("orders", new MultiSchema("orders", registry.get("orders"))); // 注册users子Schema SchemaPlus usersSchema = rootSchema.add("users", new MultiSchema("users", registry.get("users"))); } public CalciteConnection getConnection() { return calciteConnection; } }有一个细节容易踩坑:add方法的第一个参数是Schema名,这个名称会直接出现在SQL里。在做SQL编写时一定要引用准确的名称,大小写也会被Calcite严格处理。如果数据库用户名或表名带下划线还混合大写,建议在创建Schema时统一用小写,并在实际查询时给表名加双引号包裹,否则很容易出现“Table not found”的诡异报错。
3.4 执行查询并映射结果
一旦CalciteConnection建立,后续的查询就跟普通JDBC一样,获取Statement,执行SQL,遍历ResultSet。不过这里面有个小坑:Calcite自己有一套ResultSet类型转换,返回值被包装成了Calcite的游标实现,字段类型可能会跟你预想的不太一样。所以在结果映射时我统一走Object类型,由下游代码决定怎么转换。
public List<Map<String, Object>> executeQuery(String sql) throws SQLException { List<Map<String, Object>> result = new ArrayList<>(); try (Statement statement = calciteConnection.createStatement(); ResultSet rs = statement.executeQuery(sql)) { ResultSetMetaData metaData = rs.getMetaData(); int columnCount = metaData.getColumnCount(); while (rs.next()) { Map<String, Object> row = new HashMap<>(columnCount); for (int i = 1; i <= columnCount; i++) { row.put(metaData.getColumnLabel(i), rs.getObject(i)); } result.add(row); } } return result; }注意两点:一是所有真实数据库连接数相乘,因为Calcite会为每个参与执行的表拉取连接;二是不要在这里面做事务控制,Calcite不保证跨数据源的事务一致性,它只负责查询,不负责提交。这是在线报表场景下完全可以接受的前提。
4. 实操案例:一条SQL跨MySQL和PostgreSQL关联查询
先说明测试环境的情况:MySQL源orders库里有张t_order表,字段包括order_id、user_id、amount、create_time,大约五十万行数据;PostgreSQL源users库里是t_user表,字段包括user_id、user_name、email、register_time,大约十万行数据。目标是把两张表做关联,查出来一笔订单对应的用户姓名和邮箱,并且用where条件过滤掉金额为0的记录。
4.1 测试数据准备
为了让结果有区分度,我在MySQL里插入几条特定的测试订单,比如user_id分别为1001和1002的订单各两条,其中一条金额为0,另一条金额为199.99。PostgreSQL里则保证user_id 1001存在,而1002不存在。这样既能验证正常的JOIN,也能暴露多源关联时空数据处理的情况。
这个准备过程看似简单,实际上很有必要,因为联邦查询的场景下,结果依赖两个库各自的数据质量。我建议你也在测试阶段造这样一批“故意不干净”的数据,不然等联调时才发现跨源JOIN的左连接右连接语义问题,定位成本会高很多。
4.2 跨源查询SQL与执行
执行时我写的是这条SQL:
SELECT o.order_id, u.user_name, u.email, o.amount FROM orders.t_order o LEFT JOIN users.t_user u ON o.user_id = u.user_id WHERE o.amount > 0 ORDER BY o.amount DESC在Calcite里,表名的引用方式需要写成【Schema名.表名】。因为我前面注册了两个顶层Schema,所以SQL里直接写orders.t_order和users.t_user。Calcite会去对应的Schema中找到表名,然后把整条SQL拆成两部分:一份下推到MySQL执行,一份下推到PostgreSQL执行,再在Calcite内部完成JOIN。
这里有个关键点:Calcite是否把过滤条件下推,取决于它能不能推断出底层表本身也支持这个过滤表达式。像amount > 0这种普通比较条件,Calcite可以下推到MySQL节省传输量。但如果条件是函数表达式,比如DATE(create_time) = '2024-01-01',Calcite不一定能识别为可下推,这时候查性能会肉眼可见地变慢。我的建议是在自定义Table层把可下推的函数白名单扩展一下,或者直接用数据库侧能理解的普通写法。
4.3 执行结果说明与性能观察
在上面准备的测试数据下,执行结果与预期一致:amount大于0的订单返回,关联不到user_id 1002的订单时,user_name和email字段为null。这说明Calcite对LEFT JOIN的语义处理是完整的,没有出现只返回两个源都有的匹配行这种初级错误。
性能方面,五十万行的MySQL表与十万行的PostgreSQL表做JOIN,首次查询耗时约三秒,因为Calcite需要从两个源拉全量数据做内存关联。第二次相同查询降到六百毫秒左右,Calcite内部对表和Schema的元数据做了缓存。如果你的数据量动辄千万级,这个方案就不太合适了,因为Calcite的关联算法默认是内存式的HashJoin或NestedLoopJoin,数据量过大会OOM。真实生产上我更推荐用它做明细级或轻度聚合级的查询,而不是超大宽表扫描。
4.4 执行计划怎么分析
Calcite提供了一个explain功能,可以在真实执行之前查看它的执行计划,这个能力在实际排查时很有价值。通过执行EXPLAIN PLAN FOR加上你的查询SQL,能看到Calcite决定把哪些操作下推、哪些留在自己这边算。
从我的测试计划看,Calcite把对t_order的扫描和对t_user的扫描都标记为JDBCToEnumerableConverter,意味着两张表的数据都会拉到内存。这个输出其实是在提醒我:当前场景下JOIN完全发生在Calcite内存中,如果数据集膨胀,这里就是瓶颈。通过这个输出,你可以比较直观地判断某个查询能否真正下推,而不是靠猜。
5. 实战踩坑记录与问题排查
以下问题都是我在这套方案从 demo 到生产过程中真实遇到过的,按从高到低的出现频率排列。
5.1 大小写敏感与标识符符号
Calcite默认大小写策略是大小写不敏感,除非你在SQL里给标识符加了双引号。而MySQL在Linux环境下表名大小写敏感,PostgreSQL对未加引号的标识符一律折叠成小写。这个差异带来的坑是:如果同一套表名在两个库里一个大写一个小写,你在Calcite层写SQL时很容易对不上。
我的解法是统一规范,所有注册进Calcite的Schema名和表名,在构建TableMap时统一处理成小写,并在查询文档里要求SQL中使用小写表名。如果遇到本来就带大小写混合的库表,则在SQL中用双引号把表名包裹起来,比如"OrderTable",这样Calcite会严格按照引号内的名称找表。
5.2 方言识别不准导致SQL下推失败
Calcite对每个数据源都要选择合适的SQL Dialect。如果Dialect选错了,轻则优化器无法识别可下推的操作,重则生成的SQL语法不兼容,查询直接报错。我遇到过PostgreSQL方言被识别成MySQL方言的情况,布尔类型的处理方式完全不同,一个用TRUE/FALSE,一个用1/0,查出来结果经常对不上。
解决办法是注册Schema时显式指定Dialect,不要依赖Calcite自动识别。在JdbcSchema.create方法里把DatabaseProduct的对应Dialect实例传进去,比如PostgreSQL就传PostgresqlDialect.INSTANCE,MySQL传MysqlDialect.INSTANCE,代码里写死,避免运行时判断误差。
5.3 连接池耗尽与线程池配置
Calcite在跨源查询时,会并行从多个数据源拉数据。如果每个源的最大连接数都给得不大,比如默认10,而同时有十个查询任务进来,每个查询又要占两个源各一个连接,连接池很快就满了。之后新查询只能阻塞等待,表现就是接口整体卡死。
我踩到这个坑后把每个数据源的最大连接数提到50,并为Calcite查询单独开了一个FixedThreadPool,避免日常业务的数据库连接池和Calcite用同一池导致互相干扰。线程池的核心线程数按并发查询数来定,我这边核心线程数是CPU核数+1,最大线程数设了32,队列用了SynchronousQueue,宁可排队也不无限堆积任务。
5.4 元数据刷新机制
表结构变化之后,Calcite的Schema缓存不会自动失效。我在上线后加了一个管理接口,接收到刷新请求时,把对应数据源的TableMap缓存清掉并重新加载。注意清缓存的时候要保持线程安全,我用的是读写锁:刷新时加写锁,查询时加读锁。这个接口我放在内网链路里,不对外开放,同时每次刷新后打一条带耗时和表数量的日志,方便事后追溯。
5.5 常见问题速查表
| 现象 | 可能原因 | 解决方案 |
|---|---|---|
| Table not found | Schema名或表名大小写不匹配 | 统一小写,或在SQL中双引号包裹标识符 |
| SQL语法一直解析报错 | 方言选择错误,如PG被当MySQL | 注册Schema时显式指定Dialect |
| 查询结果字段全为null | 类型映射问题,如PG的timestamp类型被映射成Object | 自定义类型转换器,统一转String |
| 启动时NoSuchMethodError | guava等依赖版本冲突 | 用dependencyManagement钉版本 |
| 并发查询时接口卡住 | 连接池被占满,无空闲连接 | 调大连接池并隔离Calcite专用线程池 |
| 性能比直接查两个库还慢 | Calcite内存JOIN且无法下推 | 检查执行计划,改写SQL白名单下推条件 |
6. 这套方案的扩展边界与使用建议
最后再说一下我对这套方案边界的认识。它解决的是“多个异构数据源之间的轻量级关联查询”,最适合报表看板、运营后台、在线管理工具这类低并发但SQL灵活度要求高的场景。它替代不了数仓,也不适合大数据量的离线计算,但作为SpringBoot3应用内的“逻辑数据仓库”,能极大减少业务代码里手写内存关联的复杂度。
我后续打算在这个基础上做两件事。一是把数据源注册改成配置中心动态下发,这样新增数据源不用重启应用。二是把SQL执行链路接入一个简单的审计日志,记录每次跨源查询的SQL、耗时和下推信息,方便排查线上问题。建议你在做类似设计时也提前考虑这两点,等业务用起来再补会比较被动。
这套方案的另外一个价值是学习收益。实现过程中你会把JDBC、异构数据库方言、元数据管理、SQL解析执行原理这些零散的知识都串起来,这也是我写这篇笔记的原始动机。