SpringBoot3集成Calcite实现多数据源跨库联邦查询
开头先从实际痛点切入带出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.4JDK要求17以上这没得商量SpringBoot3强制JDK17起步。如果你还在用JDK8那第一件事就是先把运行时升上去。项目本身就是一个普通的SpringBoot Web工程我习惯先建一个干净的骨架不要急着堆依赖越干净后面排错越容易。Calcite这边我用的版本是1.36.0在Maven中央仓库直接能拉。要引入的核心依赖是两个calcite-core和calcite-linq4j。前者提供SQL解析、校验、优化和核心的Schema模型后者是Calcite内部做表达式计算和数据集遍历用的很多教程只提core结果运行时报ClassNotFoundException这里直接放在pom里。dependency groupIdorg.apache.calcite/groupId artifactIdcalcite-core/artifactId version1.36.0/version /dependency dependency groupIdorg.apache.calcite/groupId artifactIdcalcite-linq4j/artifactId version1.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 MapString, 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 CollectionString 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 groupIdcom.google.guava/groupId artifactIdguava/artifactId version32.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 MapString, org.apache.calcite.schema.Table getTableMap() { // 从数据源获取元数据构建表名到JdbcTable的映射 return buildTableMap(); } }getTableMap的构建逻辑重点是拿到真实数据库里所有物理表的表名然后为每个表名创建JdbcTable。JdbcTable的构造需要传入一个JdbcSchema实例和表名。如果数据库表非常多这一步会比较重所以我加了一层缓存应用启动时初始化一次之后通过刷新接口在表结构变更时手动触发重建。3.3 在SpringBoot中注册Schema并获取CalciteConnectionSchema构建好之后需要把它挂到一个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 ListMapString, Object executeQuery(String sql) throws SQLException { ListMapString, 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()) { MapString, 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与执行执行时我写的是这条SQLSELECT 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-01Calcite不一定能识别为可下推这时候查性能会肉眼可见地变慢。我的建议是在自定义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.INSTANCEMySQL传MysqlDialect.INSTANCE代码里写死避免运行时判断误差。5.3 连接池耗尽与线程池配置Calcite在跨源查询时会并行从多个数据源拉数据。如果每个源的最大连接数都给得不大比如默认10而同时有十个查询任务进来每个查询又要占两个源各一个连接连接池很快就满了。之后新查询只能阻塞等待表现就是接口整体卡死。我踩到这个坑后把每个数据源的最大连接数提到50并为Calcite查询单独开了一个FixedThreadPool避免日常业务的数据库连接池和Calcite用同一池导致互相干扰。线程池的核心线程数按并发查询数来定我这边核心线程数是CPU核数1最大线程数设了32队列用了SynchronousQueue宁可排队也不无限堆积任务。5.4 元数据刷新机制表结构变化之后Calcite的Schema缓存不会自动失效。我在上线后加了一个管理接口接收到刷新请求时把对应数据源的TableMap缓存清掉并重新加载。注意清缓存的时候要保持线程安全我用的是读写锁刷新时加写锁查询时加读锁。这个接口我放在内网链路里不对外开放同时每次刷新后打一条带耗时和表数量的日志方便事后追溯。5.5 常见问题速查表现象可能原因解决方案Table not foundSchema名或表名大小写不匹配统一小写或在SQL中双引号包裹标识符SQL语法一直解析报错方言选择错误如PG被当MySQL注册Schema时显式指定Dialect查询结果字段全为null类型映射问题如PG的timestamp类型被映射成Object自定义类型转换器统一转String启动时NoSuchMethodErrorguava等依赖版本冲突用dependencyManagement钉版本并发查询时接口卡住连接池被占满无空闲连接调大连接池并隔离Calcite专用线程池性能比直接查两个库还慢Calcite内存JOIN且无法下推检查执行计划改写SQL白名单下推条件6. 这套方案的扩展边界与使用建议最后再说一下我对这套方案边界的认识。它解决的是“多个异构数据源之间的轻量级关联查询”最适合报表看板、运营后台、在线管理工具这类低并发但SQL灵活度要求高的场景。它替代不了数仓也不适合大数据量的离线计算但作为SpringBoot3应用内的“逻辑数据仓库”能极大减少业务代码里手写内存关联的复杂度。我后续打算在这个基础上做两件事。一是把数据源注册改成配置中心动态下发这样新增数据源不用重启应用。二是把SQL执行链路接入一个简单的审计日志记录每次跨源查询的SQL、耗时和下推信息方便排查线上问题。建议你在做类似设计时也提前考虑这两点等业务用起来再补会比较被动。这套方案的另外一个价值是学习收益。实现过程中你会把JDBC、异构数据库方言、元数据管理、SQL解析执行原理这些零散的知识都串起来这也是我写这篇笔记的原始动机。