You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.21 09:48:16