ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

SpringBoot3集成Apache Calcite:多数据源跨库查询实战指南

SpringBoot3集成Apache Calcite:多数据源跨库查询实战指南 前一阵做报表中心的时候我接到一个很拧巴的需求一张查询页里要同时带出订单库、用户库和库存库的数据筛选条件还横跨三个库数据要实时不能走同步任务。最开始按SpringBoot多数据源的思路走结果业务代码越写越别扭——查完订单表还要拿userId去用户库批查最后在Java内存里做关联这不是在写SQL是在自己写SQL引擎。后来我改用Apache Calcite花了一周把它集成到SpringBoot3里把多数据源查询变成了一条简单的SQL。这篇就是当时的实战笔记重点讲SpringBoot3集成Calcite做多数据源查询的完整思路、关键代码和踩坑记录适合正在折腾多源查询、异构库join的后端同学参考。1. 多源查询的老路走不通为什么我把目光转向Calcite1.1 一张报表背后牵出三个库我们的业务不算复杂订单表在MySQL订单库用户信息在另一个MySQL用户库库存和商品快照在PostgreSQL。报表需求是查出“最近30天订单金额超过500元且用户所在城市为杭州同时附带商品库存变化”。用SQL描述其实很简单就是三张表关联加过滤。但表不在同一个库里问题就变成了“怎么让几个库像一张库一样被查”。这应该是很多团队都会遇到的多数据源查询场景。数据量不算海量几十万到百万级但对实时性有要求不能接受昨晚的T1数据。开发同学也不想为了一个查询去写三套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.annotationSpringBoot3的Jakarta注解体系和它容易冲突启动时会出现奇怪的注解扫描异常不建议在SpringBoot3新项目里用老版本。Maven依赖很简单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、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(ListVirtualSchema 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的对象模型是一棵树。最顶层是RootSchemaRootSchema下面可以挂一个或多个自定义Schema每个自定义Schema里再挂Table。查询时用的orders.oms_order意思是RootSchema下面有个orders的Schema这个Schema里有一张oms_order的表。这个模型和我的多数据源需求正好对应。我建了三个逻辑Schemaorders对应MySQL订单库user对应MySQL用户库stock对应PostgreSQL库存库。每个Schema内部维护了一个表名到Table实现的映射。物理数据源本身还是用Spring管理的DataSourceCalcite不接管连接池只通过抽象接口去访问物理表。这个设计把SpringBoot生态和Calcite的查询引擎解耦得很干净两边各管各的。3. 自定义Schema与Table把MySQL/PostgreSQL表暴露成虚拟表3.1 让Calcite认识MySQL表AbstractTable的scan方法注册Schema比较简单重写AbstractSchema的getTableMap()就行public class DynamicSchema extends AbstractSchema { private final MapString, Table tableMap; public DynamicSchema(MapString, Table tableMap) { this.tableMap tableMap; } Override protected MapString, 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 EnumerableObject[] scan(DataContext root) { String sql select * from tableName; return new AbstractEnumerable() { Override public EnumeratorObject[] 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.DECIMALPostgreSQL的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之后复杂查询就变成了一条普通SQLSELECT 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方法会额外收到一个ListRexNode filters里面就是SQL中下推到这一层的过滤条件。我的做法是做一个PushdownTable extends AbstractTable implements FilterableTable在scan里把filters翻译成目标数据库的WHERE条件拼到查询SQL里。简单场景下的翻译逻辑可以这样写Override public EnumerableObject[] scan(DataContext root, ListRexNode filters) { String whereSql ; if (!filters.isEmpty()) { ListString 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设置为JAVACalcite对大小写敏感两个名字不匹配。顺带发现另一个坑如果SQL里没有带Schema前缀比如直接写select * from oms_orderCalcite会在默认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 10)) { 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 10查元数据而不是select * from table是为了避免全表扫描只拿表结构性能很好。这也是我想要提醒的不要让类型映射靠手写让JDBC元数据告诉你真相。5.3 连接生命周期与并发安全Calcite的执行是拉取式的这带来一个容易被忽略的问题连接关闭的职责落在了Enumerator的close()上。如果scan里从DataSource拿了物理连接一定要在Enumerator关闭时释放。我当时封装JdbcEnumerator在落SQL执行后的Enumerator.close里统一释放ResultSet、PreparedStatement和ConnectionOverride 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真正放到生产环境还有一些工程层面的加固动作。我整理了自己落地时的几项必做配置限制只读SQLCalcite连接默认暴露的接口能力较完整不能把任意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。希望这篇实战笔记能帮你少走几个弯路。
返回列表