如何在Apache Calcite中注册可下推至数据源的UDF
问题描述
我已经注册了一个自定义UDFformat_datetime,用于将Timestamp转换为指定格式的字符串,代码如下:
public static class ConvertTimestampToString { public final static String NAME = "format_datetime"; private static final Map<String, DateTimeFormatter> FORMATTERS = new WeakHashMap<>(); public final static ScalarFunction FUNCTION = ScalarFunctionImpl.create( Types.lookupMethod( ConvertTimestampToString.class, "eval", Timestamp.class, String.class)); public final static SqlIdentifier IDENTIFIER = new SqlIdentifier( Collections.singletonList(NAME), null, SqlParserPos.ZERO, null); public final static SqlOperator SQL_OPERATOR = new SqlUserDefinedFunction( IDENTIFIER, ReturnTypes.VARCHAR_2000, null, OperandTypes.family(SqlTypeFamily.TIMESTAMP, SqlTypeFamily.STRING), null, FUNCTION); public String eval(Timestamp source, String format) { try { var formatter = FORMATTERS .computeIfAbsent(format, f -> DateTimeFormatter.ofPattern(f)); return source.toLocalDateTime().format(formatter); } catch (Exception e) { return null; } } }
使用RelBuilder构建查询时,代码如下:
RelNode node = builder .scan(source) .filter(builder.isNotNull(builder.field("valid_to_dttm"))) .project(builder.call( ConvertTimestampToString.SQL_OPERATOR, builder.field("valid_to_dttm"), builder.literal("YYYY/MM"))) .limit(0, 10) .build(); try (PreparedStatement st = relRunner.prepareStatement(node); ResultSet rs = st.executeQuery()) { QueryLogsUtil.logTrinoSqlQuery(rs); printResultSet(rs); }
但当前生成的SQL并未将format_datetime推送到数据源,而是在Calcite本地计算:
LogicalSort(fetch=[10]) LogicalProject(resolution_dttm=[format_datetime($0, _UTF-16'MM/YYYY':CHAR(7) CHARACTER SET "UTF-16")]) LogicalFilter(condition=[IS NOT NULL($0)]) JdbcTableScan(table=[[trino, m1-gp-tst-ext, hercule_test, jira_issue]]) SELECT * FROM "m1-gp-tst-ext"."hercule_test"."jira_issue" WHERE "resolution_dttm" IS NOT NULL FETCH NEXT 10 ROWS ONLY
我需要调整配置,让生成的SQL变为SELECT format_datetime(valid_to_dttm, 'YYYY/MM') FROM ...,即把UDF下推到数据源执行。
解决方案
要让Calcite将UDF下推到JDBC数据源,需要完成以下3个关键步骤:
1. 标记UDF为可下推
修改SqlUserDefinedFunction的定义,添加函数类别和可下推特性,告诉Calcite这个函数可以被下推到数据源:
public final static SqlOperator SQL_OPERATOR = new SqlUserDefinedFunction( IDENTIFIER, ReturnTypes.VARCHAR_2000, null, OperandTypes.family(SqlTypeFamily.TIMESTAMP, SqlTypeFamily.STRING), Collections.singletonList(SqlFunctionCategory.SYSTEM), // 指定函数类别,表明是数据源支持的函数 FUNCTION) { @Override public boolean isDeterministic() { return true; // 标记函数为确定性(相同输入返回相同输出),这是下推的前提 } @Override public boolean canImplement(SqlToRelConverter converter) { return false; // 禁止Calcite在本地实现该函数,强制下推到数据源 } };
2. 配置JDBC方言支持该UDF
如果使用的是Trino JDBC,需要确保Calcite的Trino方言识别format_datetime函数。可以自定义方言,将该函数添加到支持的算子列表:
public class CustomTrinoDialect extends JdbcDialect { public CustomTrinoDialect(JdbcConvention convention) { super(convention); } @Override public boolean supportsFunction(SqlOperator operator) { if (operator instanceof SqlUserDefinedFunction) { return ConvertTimestampToString.NAME.equals(((SqlUserDefinedFunction) operator).getNameAsId().getSimple()); } return super.supportsFunction(operator); } @Override public void unparseCall(SqlWriter writer, SqlCall call, int leftPrec, int rightPrec) { if (ConvertTimestampToString.SQL_OPERATOR.equals(call.getOperator())) { // 按照Trino的语法格式输出函数调用 writer.keyword(ConvertTimestampToString.NAME); writer.print("("); call.operand(0).unparse(writer, leftPrec, rightPrec); writer.print(", "); call.operand(1).unparse(writer, leftPrec, rightPrec); writer.print(")"); } else { super.unparseCall(writer, call, leftPrec, rightPrec); } } }
然后在创建JdbcConvention时使用自定义方言:
JdbcConvention convention = JdbcConvention.of(datasource, CustomTrinoDialect::new);
3. 启用Calcite的下推优化规则
确保Calcite的优化器启用了JdbcProjectRule和JdbcFilterRule等下推规则。在创建RelOptPlanner时,添加这些规则:
RelOptPlanner planner = Frameworks.getPlanner(config); planner.addRule(JdbcRules.JdbcProjectRule.INSTANCE); planner.addRule(JdbcRules.JdbcFilterRule.INSTANCE); planner.addRule(JdbcRules.JdbcSortRule.INSTANCE);
或者在RelBuilder的配置中启用优化:
RelBuilder builder = RelBuilder.create(config); builder.getPlanner().addRule(JdbcRules.JdbcProjectRule.INSTANCE);
验证调整后的结果
完成上述配置后,重新构建查询,生成的SQL应该变为:
SELECT format_datetime("valid_to_dttm", 'YYYY/MM') FROM "m1-gp-tst-ext"."hercule_test"."jira_issue" WHERE "valid_to_dttm" IS NOT NULL FETCH NEXT 10 ROWS ONLY
对应的RelNode计划会显示JdbcProject直接在JdbcTableScan之上,说明UDF已被下推到数据源执行。
内容的提问来源于stack exchange,提问作者kalmar
相关产品推荐
相关产品推荐

