前一阵做报表中心的时候,我接到一个很拧巴的需求:一张查询页里要同时带出订单库、用户库和库存库的数据,筛选条件还横跨三个库,数据要实时,不能走同步任务。最开始按SpringBoot多数据源的思路走,结果业务代码越写越别扭——查完订单表,还要拿userId去用户库批查,最后在Java内存里做关联,这不是在写SQL,是在自己写SQL引擎。后来我改用Apache Calcite,花了一周把它集成到SpringBoot3里,把多数据源查询变成了一条简单的SQL。这篇就是当时的实战笔记,重点讲SpringBoot3集成Calcite做多数据源查询的完整思路、关键代码和踩坑记录,适合正在折腾多源查询、异构库join的后端同学参考。
1. 多源查询的老路走不通:为什么我把目光转向Calcite
1.1 一张报表背后牵出三个库
我们的业务不算复杂:订单表在MySQL订单库,用户信息在另一个MySQL用户库,库存和商品快照在PostgreSQL。报表需求是查出“最近30天订单金额超过500元,且用户所在城市为杭州,同时附带商品库存变化”。用SQL描述其实很简单,就是三张表关联加过滤。但表不在同一个库里,问题就变成了“怎么让几个库像一张库一样被查”。
这应该是很多团队都会遇到的多数据源查询场景。数据量不算海量,几十万到百万级,但对实时性有要求,不能接受昨晚的T+1数据。开发同学也不想为了一个查询去写三套Mapper接口再在Service层手动for循环拼结果。
1.2 常规多源方案的三个坎
我把团队里常用的方案拉出来比了一圈,每个都有明显瓶颈。
第一种是配置多个DataSource,把多个库的Mapper分目录管理。这个方案最适合“一个事务里只碰一个库”的场景,但碰到跨库join就抓瞎:要么在Java里分页拉出来再手工合并,要么写临时表把数据倒来倒去。代码里到处都是List循环、Map关联,维护成本很高。
第二种是引入数据库中间件,走分库分表路由。能力很强,SQL路由、读写分离、分布式事务都有,但对一个原本没有分库分表需求的系统来说,部署和配置成本偏高,很多中间件还会拦截你的SQL做改写,一旦碰到自定义函数、复杂子查询,排查问题会非常痛苦。
第三种是做ETL同步,把多个库的数据汇聚到数仓或宽表。这是最“稳定”的,但也是实时性最差的。为了一个临时报表需求专门建同步链路,运维同学想打人。同步任务的延迟、失败重试、字段变更都要额外处理。
1.3 为什么要选Calcite:它把自己定位成什么
Apache Calcite被很多同学误解成搜索引擎或者分布式查询引擎,这些都不是。它的本质是一个SQL处理框架:负责SQL的解析、校验、逻辑优化、物理优化,再通过可插拔的Adapter去访问外部数据源。它不存储数据,只负责“听懂SQL,然后指挥各个数据源干活”。
这个定位对多数据源查询来说太合适了。我不需要把数据搬到一个地方,只需要让Calcite知道每个数据源有哪些库、哪些表、字段类型是什么。查询时它把SQL拆成逻辑计划,能从物理库下推的就下推,不能下推的比如跨库join,它才在内存里替我做。按我之前做的对比:
| 对比维度 | 多DataSource MyBatis | 数据库中间件 | ETL同步 | Calcite嵌入式 |
|---|---|---|---|---|
| 统一SQL能力 | 不支持 | 支持 | 不支持 | 支持 |
| 跨库join | 需手写聚合 | 支持 | 支持 | 支持 |
| 代码侵入性 | 中 | 中高 | 高 | 中 |
| 部署依赖 | 无 | 需额外节点 | 需同步集群 | 纯依赖引入 |
| 实时性 | 原库实时 | 原库实时 | 有延迟 | 原库实时 |
| 学习成本 | 低 | 偏高 | 偏高 | 中 |
最终我选了Calcite,把报表中心的核心查询能力搭在它上面。下面从项目落地角度,把步骤和坑都记录下来。
2. 环境与版本:SpringBoot3项目集成Calcite的依赖和启动模型
2.1 版本兼容性:SpringBoot3与Calcite的配平
先讲版本,这是SpringBoot3集成Calcite最容易翻车的地方。SpringBoot3要求JDK17起步,底层是Jakarta EE规范。Calcite这边,我实测到1.36.0是稳定可用的。1.33及以前的版本还在用javax.annotation,SpringBoot3的Jakarta注解体系和它容易冲突,启动时会出现奇怪的注解扫描异常,不建议在SpringBoot3新项目里用老版本。
Maven依赖很简单:
<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、Avatica等。SpringBoot的dependencyManagement会统一管Jackson的版本,两边的版本如果不一致,轻则日志警告,重则反序列化直接报错。我当时的做法是引完依赖后跑一次mvn dependency:tree,把Calcite里传递进来的旧版Guava用exclusion排除掉,让SpringBoot的版本兜底。这一步没什么技术含量,但漏做的同学后面会踩得很惨。
2.2 一段能跑起来的最小引擎代码
SpringBoot3里集成Calcite,不需要额外中间件,核心就是拿到一个CalciteConnection。这个连接对象是java.sql.Connection的实现类,通过JDBC URLjdbc:calcite:创建。我封装了一个最小的引擎类:
public class CalciteEngine { private final CalciteConnection connection; public CalciteEngine(List<VirtualSchema> schemas) throws SQLException { Properties info = new Properties(); info.setProperty("lex", "JAVA"); this.connection = DriverManager.getConnection("jdbc:calcite:", info) .unwrap(CalciteConnection.class); SchemaPlus rootSchema = connection.getRootSchema(); for (VirtualSchema schema : schemas) { rootSchema.add(schema.getName(), schema.toCalciteSchema()); } } public void query(String sql) throws SQLException { try (Statement statement = connection.createStatement(); ResultSet rs = statement.executeQuery(sql)) { while (rs.next()) { // 按列取值处理 } } } }这段代码里有两个关键点。第一个是info.setProperty("lex", "JAVA"),它决定了Calcite对SQL标识符大小写的处理策略。选JAVA时,Calcite会按Java标识符的规则处理大小写,SQL里的表名和字段名必须和注册时完全一致。第二个是拿到CalciteConnection后,必须通过getRootSchema()往根Schema上挂数据源Schema,后续SQL才能用schemaName.tableName的方式访问。
2.3 从RootSchema理解Calcite的对象模型
Calcite的对象模型是一棵树。最顶层是RootSchema,RootSchema下面可以挂一个或多个自定义Schema,每个自定义Schema里再挂Table。查询时用的orders.oms_order,意思是RootSchema下面有个orders的Schema,这个Schema里有一张oms_order的表。
这个模型和我的多数据源需求正好对应。我建了三个逻辑Schema:orders对应MySQL订单库,user对应MySQL用户库,stock对应PostgreSQL库存库。每个Schema内部维护了一个表名到Table实现的映射。物理数据源本身还是用Spring管理的DataSource,Calcite不接管连接池,只通过抽象接口去访问物理表。这个设计把SpringBoot生态和Calcite的查询引擎解耦得很干净,两边各管各的。
3. 自定义Schema与Table:把MySQL/PostgreSQL表暴露成虚拟表
3.1 让Calcite认识MySQL表:AbstractTable的scan方法
注册Schema比较简单,重写AbstractSchema的getTableMap()就行:
public class DynamicSchema extends AbstractSchema { private final Map<String, Table> tableMap; public DynamicSchema(Map<String, Table> tableMap) { this.tableMap = tableMap; } @Override protected Map<String, Table> getTableMap() { return tableMap; } }真正麻烦的是Table实现。Calcite暴露给外部数据源的访问接口有好几个:ScannableTable、FilterableTable、TranslatableTable。最小可用的是AbstractTable加scan方法,官方CSV例子就是这么干的。我写了一个PhysicalTable,把Spring的DataSource注入进来,scan时真正去物理库执行查询:
public class PhysicalTable extends AbstractTable { private final DataSource dataSource; private final String tableName; private final RelDataType rowType; public PhysicalTable(DataSource dataSource, String tableName, RelDataType rowType) { this.dataSource = dataSource; this.tableName = tableName; this.rowType = rowType; } @Override public RelDataType getRowType(RelDataTypeFactory typeFactory) { return rowType; } @Override public Enumerable<Object[]> scan(DataContext root) { String sql = "select * from " + tableName; return new AbstractEnumerable<>() { @Override public Enumerator<Object[]> enumerator() { try { Connection conn = dataSource.getConnection(); PreparedStatement ps = conn.prepareStatement(sql); ResultSet rs = ps.executeQuery(); return new JdbcEnumerator(rs, conn, ps); } catch (SQLException ex) { throw new RuntimeException("查询物理表失败: " + tableName, ex); } } }; } }Calcite的执行模型是拉取式的。它不会一次性把结果集全读进来,而是拿到我返回的Enumerator之后,由上层算子一个个moveNext()去拉数据。这意味着我在scan里不应该把整个ResultSet转成List再返回,直接把ResultSet包成Enumerator反馈给Calcite才是最省内存的做法。
3.2 数据类型映射:这是最容易翻车的地方
构造PhysicalTable时,必须提供RelDataType,这是Calcite对表结构的“认知”。它决定了SQL表达式计算、字段类型转换、结果集类型强转的基准。如果这里给错了,后面会报一些非常抽象的异常。
我一开始图省事,把所有字段都映射成VARCHAR,结果SQL里写where amount > 100时,Calcite在把字符串和数字比较时直接抛异常。后来老老实实按物理表元数据构建:
RelDataTypeFactory typeFactory = new SqlTypeFactoryImpl(RelDataTypeSystem.DEFAULT); RelDataTypeFactory.Builder builder = typeFactory.builder(); builder.add("order_id", typeFactory.createSqlType(SqlTypeName.BIGINT)); builder.add("user_id", typeFactory.createSqlType(SqlTypeName.BIGINT)); builder.add("amount", typeFactory.createSqlType(SqlTypeName.DECIMAL)); builder.add("status", typeFactory.createSqlType(SqlTypeName.INTEGER)); RelDataType rowType = builder.build();这里的原则是:逻辑表字段类型要和物理表JDBC类型尽量对齐,不能想当然。MySQL的DECIMAL精确对应SqlTypeName.DECIMAL,PostgreSQL的NUMERIC也一样。你在Spring里怎么用BigDecimal接,Calcite这里就怎么声明。
3.3 懒加载机制:为什么我的表“找不到”
Calcite有个很容易误导人的特性:完全懒加载。rootSchema.add()并不会去连接物理库,也不会校验表是否存在、字段是否匹配。只有真正执行SQL时,Calcite才会调用对应Table的getRowType和scan。好处是启动速度快,坏处是表名拼错、字段写错只能在查询运行时暴雷,不会在配置阶段提前报警。
我在调试阶段就遇到过反复报“Table not found”的情况,跟这个特性强相关。具体排查过程我放到后面第五部分细讲,先记住一个结论:注册Schema时不要试图做任何远程校验,校验逻辑放到查询入口统一做。你可以在应用启动后用一条select count(*) from xx.yy做冒烟测试,但这个测试本身就是第一次真实查询。
4. 跨库Join与过滤条件下推:Calcite到底是怎么干活的
4.1 一条跨库Join SQL在Calcite里走了什么路
注册好三个Schema之后,复杂查询就变成了一条普通SQL:
SELECT o.order_id, u.user_name, s.stock_qty FROM orders.oms_order o JOIN user.uc_user u ON o.user_id = u.id LEFT JOIN stock.goods_stock s ON o.goods_id = s.goods_id WHERE o.status = 1 AND u.city = '杭州'这条SQL的执行路径大致是:Calcite的Parser把文本转成语法树,Validator做表名、字段名、类型检查,之后转换成关系代数,再经过优化器生成物理执行计划。默认模式下,Calcite会把oms_order、uc_user、goods_stock三张表分别拉取出来,在内存中做HashJoin或MergeJoin,最后过滤输出。
这意味着什么?意味着如果你的过滤条件没下推,Calcite会先把整个oms_order表全量拉到内存,再全量拉uc_user,再全量拉goods_stock,最后三张表在JVM堆里join。表小没事,表一大直接把自己玩死。
4.2 FilterableTable:过滤条件下推的落地姿势
要让过滤条件尽量在物理库执行,就不能只实现AbstractTable,要实现FilterableTable接口。这个接口的scan方法会额外收到一个List<RexNode> filters,里面就是SQL中下推到这一层的过滤条件。
我的做法是做一个PushdownTable extends AbstractTable implements FilterableTable,在scan里把filters翻译成目标数据库的WHERE条件,拼到查询SQL里。简单场景下的翻译逻辑可以这样写:
@Override public Enumerable<Object[]> scan(DataContext root, List<RexNode> filters) { String whereSql = ""; if (!filters.isEmpty()) { List<String> conditions = new ArrayList<>(); for (RexNode filter : filters) { if (filter instanceof RexCall call) { SqlKind kind = call.getOperator().getKind(); if (kind == SqlKind.EQUALS) { RexInputRef ref = (RexInputRef) call.getOperands().get(0); RexLiteral literal = (RexLiteral) call.getOperands().get(1); String columnName = rowType.getFieldList().get(ref.getIndex()).getName(); conditions.add(columnName + " = " + literal.getValueAs(String.class)); } } } whereSql = " where " + String.join(" and ", conditions); } String sql = "select * from " + tableName + whereSql; // 继续走外层拉取逻辑 }这里有个细节:RexInputRef的索引是Calcite逻辑计划里的项目索引,不是物理表里的第几列。所以翻译条件前要通过rowType.getFieldList().get(ref.getIndex()).getName()拿到真正的列名,再拼SQL。我见过不少同学拿到RexInputRef直接用索引去物理列取数据,拼出来的条件张冠李戴,查出来的结果全错。
如果数据源方言复杂,比如有PostgreSQL的JSON操作、数组函数、MySQL的DATE_FORMAT,自己解析RexNode会非常累。更工程化的方案是用RelToSqlConverter配合SqlDialect,让Calcite把整个Filter节点翻译成目标方言SQL。这样通用性更强,但需要处理SqlDialect对函数名的映射问题,适合把下推能力做成通用能力的场景。
4.3 实测:什么样的查询会让内存爆炸
我做过一次对比测试。同一张约60万行的订单表,不做任何下推时,Calcite从MySQL拉回60万行到JVM,再和20万用户表做内存join,大概吃掉2GB堆内存,Young GC频繁,整个查询耗时8秒多。加上status = 1和city = '杭州'两个条件下推之后,订单表返回行数降到3万,用户表降到几千,查询耗时直接降到300毫秒左右,内存占用忽略不计。
所以我的经验是:小表、字典表可以做内存join,大表过滤条件必须下推。如果一张底层大表无论如何都要全量参与join,那Calcite默认执行引擎就不合适了。后续扩展方向是把Table改成TranslatableTable,实现toRel(),直接把它翻译成目标库的子查询,让数据库端自己完成join。这条路更复杂,但对大数据量场景是必要的。
5. 踩坑记录与生产加固:几个让我熬夜的真实问题
5.1 Table not found的背后是Schema懒加载
第一个让我熬夜的问题就是“Table 'OMS_ORDER' not found”。表名明明注册了,SQL也拼对了,就是找不到。
排查链路是这样的:先看rootSchema.getSubSchemaNames()有没有orders,发现Schema在;再看orders.getTableNames()有没有oms_order,也在。问题出在大小写。我注册表名时用的是小写oms_order,但SQL里写的是OMS_ORDER。由于lex设置为JAVA,Calcite对大小写敏感,两个名字不匹配。
顺带发现另一个坑:如果SQL里没有带Schema前缀,比如直接写select * from oms_order,Calcite会在默认Schema里找,而默认Schema并不是我注册的某个业务Schema。这个问题尤其隐蔽,因为单库调试时表名能命中,一旦换成多Schema环境就报not found。最终的规范是:所有生产SQL必须强制写schema.table,不要省略前缀。
5.2 Decimal/BigDecimal和类型强转的拉锯战
另一个让我头疼的问题是ClassCastException: java.math.BigDecimal cannot be cast to java.lang.Double。现象在sum(o.amount)这类聚合SQL上必现,单行查询却不报错。
根因还是rowType声明和物理类型不一致。我在某张表上把amount声明成了DOUBLE,但MySQL的DECIMAL通过JDBC取出来是BigDecimal。Calcite执行聚合时按DOUBLE的规则去强转物理值,自然就炸了。修法不是去改物理表,而是让rowType和物理表元数据对齐。
批量建表时我写了一个基于ResultSetMetaData的反向生成逻辑:
try (Connection conn = dataSource.getConnection(); Statement st = conn.createStatement(); ResultSet rs = st.executeQuery("select * from " + tableName + " where 1=0")) { ResultSetMetaData meta = rs.getMetaData(); RelDataTypeFactory.Builder builder = typeFactory.builder(); for (int i = 1; i <= meta.getColumnCount(); i++) { String columnName = meta.getColumnName(); int jdbcType = meta.getColumnType(); switch (jdbcType) { case Types.BIGINT -> builder.add(columnName, SqlTypeName.BIGINT); case Types.DECIMAL, Types.NUMERIC -> builder.add(columnName, SqlTypeName.DECIMAL); case Types.INTEGER -> builder.add(columnName, SqlTypeName.INTEGER); case Types.VARCHAR -> builder.add(columnName, SqlTypeName.VARCHAR); case Types.TIMESTAMP -> builder.add(columnName, SqlTypeName.TIMESTAMP); default -> builder.add(columnName, SqlTypeName.ANY); } } return builder.build(); }用where 1=0查元数据而不是select * from table,是为了避免全表扫描,只拿表结构,性能很好。这也是我想要提醒的:不要让类型映射靠手写,让JDBC元数据告诉你真相。
5.3 连接生命周期与并发安全
Calcite的执行是拉取式的,这带来一个容易被忽略的问题:连接关闭的职责落在了Enumerator的close()上。如果scan里从DataSource拿了物理连接,一定要在Enumerator关闭时释放。我当时封装JdbcEnumerator,在落SQL执行后的Enumerator.close里统一释放ResultSet、PreparedStatement和Connection:
@Override public void close() { try { if (rs != null) rs.close(); if (ps != null) ps.close(); if (conn != null) conn.close(); } catch (SQLException ignored) { // 关闭阶段的异常不能影响主流程 } }连接泄漏和并发问题往往是连锁反应。Calcite的CalciteConnection不是设计成多线程共享执行SQL的,我在项目里实测过,多个线程用同一个Connection并发查询时,会出现结果串行、状态错乱的情况。解决方式是建一个小型的CalciteConnection池,每个查询从池里取一个独立连接,用完归还,或者用ThreadLocal为每个线程维护独立连接。生产环境中我推荐后者:查询引擎本身无状态,但连接对象有执行状态,隔离得越彻底越省心。
5.4 生产环境建议清单
把Calcite真正放到生产环境,还有一些工程层面的加固动作。我整理了自己落地时的几项必做配置:
- 限制只读SQL:Calcite连接默认暴露的接口能力较完整,不能把任意SQL直接开放给业务方。我在查询入口先解析SQL,只放行
SELECT和EXPLAIN,遇到INSERT、UPDATE、DELETE、DDL直接拒绝。 - 设置物理连接超时:HikariCP的
connectionTimeout、JDBC URL里的socketTimeout都配上,避免一个慢库拖死整个查询线程。 - 查询超时熔断:给Statement设置
queryTimeout,这个时间要根据实际压测结果定,我默认给10秒。 - 缓存热点查询:对同一SQL同一参数的结果做5分钟短缓存,能显著降低物理库压力。
- 埋点监控:对每个Table的scan方法加耗时统计,超过500ms的查询单独捞出来分析,重点看过滤条件下推有没有生效。
生产环境里最怕的不是Calcite本身出问题,而是物理数据源的慢查询被Calcite放大。因为一条SQL可能同时拉多张表,只要其中一张表没有下推条件,整条查询的性能就会被打回原形。
我现在的做法是把Calcite查询引擎独立封装成一个模块,外面再包一层REST接口。业务方传SQL过来,我负责解析、只读校验、超时控制、结果集分页返回。这套结构稳定跑了几个月,中途最大的改动也只是把几张核心大表的ScannableTable换成了FilterableTable。如果你也在SpringBoot3里做多数据源查询,Calcite值得认真试一试,但一定要先想清楚哪些条件能下推,哪些表注定要内存join。希望这篇实战笔记能帮你少走几个弯路。