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

使用自定义UDF转换字段配置Flink水位线时程序编译失败

问题:自定义时间转换UDF生成的字段配置水位线时编译失败

因Flink Table API原生TO_TIMESTAMP(value, format)不支持yyyyMMddHHmmssSSS格式的解析,我们实现了TO_LOCALDATETIME自定义UDF完成转换。但将该UDF生成的字段配置为水位线事件时间时,程序触发编译错误,错误栈如下:

Caused by: org.apache.flink.api.common.InvalidProgramException: Table program cannot be compiled. This is a bug. Please file an issue.  at org.apache.flink.table.runtime.generated.CompileUtils.doCompile(CompileUtils.java:107) ~[flink-table-runtime-1.17.1.jar:1.17.1]  at org.apache.flink.table.runtime.generated.CompileUtils.lambda$compile$0(CompileUtils.java:92) ~[flink-table-runtime-1.17.1.jar:1.17.1] .... 

Caused by: org.codehaus.commons.compiler.CompileException: Line 31, Column 83: Cannot determine simple type name "org"  at org.codehaus.janino.UnitCompiler.compileError(UnitCompiler.java:12211) ~[flink-table-runtime-1.17.1.jar:1.17.1]  at org.codehaus.janino.UnitCompiler.getReferenceType(UnitCompiler.java:6833) ~[flink-table-runtime-1.17.1.jar:1.17.1] ....

java.lang.RuntimeException: Could not instantiate generated class 'WatermarkGenerator$0'
    at org.apache.flink.table.runtime.generated.GeneratedClass.newInstance(GeneratedClass.java:74) ~[flink-table-runtime-1.17.1.jar:1.17.1]   

该问题在Flink 1.16.1和1.17.1版本中均能复现。另外我们的TO_BIGDECIMAL自定义UDF可正常使用,且不配置水位线时,TO_LOCALDATETIME UDF也能正常工作。


相关代码

UDF定义

public class ToLocalDateTimeFunction extends ScalarFunction {
    private static final DateTimeFormatter format = new DateTimeFormatterBuilder()
            .appendPattern("yyyyMMddHHmmss")
            .appendValue(ChronoField.MILLI_OF_SECOND, 3)
            .toFormatter();
    public LocalDateTime eval(String value) {
        try {
            return LocalDateTime.parse(value, format);
        } catch (Exception e) {
            return null;
        }
    }
}

UDF注册

tableEnvironment.createTemporarySystemFunction("TO_BIGDECIMAL", ToBigDecimalFunction.class);
tableEnvironment.createTemporarySystemFunction("TO_LOCALDATETIME", ToLocalDateTimeFunction.class);
tableEnvironment.createTemporaryTable("transactionTable", transactionDataStream);

TableDescriptor中使用UDF

TableDescriptor.forConnector(KafkaDynamicTableFactory.IDENTIFIER)
        .schema(Schema.newBuilder()
                .columnByExpression("entityId", "raw_data.entity_id")
                .columnByExpression("tranDate", "TO_LOCALDATETIME(raw_data.tran_date)")
                .columnByExpression("rowtime", "CAST(TO_LOCALDATETIME(raw_data.tran_date) AS TIMESTAMP_LTZ(3))")
                .columnByExpression("amt", "TO_BIGDECIMAL(raw_data.amt)")
                .watermark("rowtime", "rowtime - INTERVAL '10' MINUTE")
                .columnByExpression("processTime", "PROCTIME()")
                .column("raw_data", DataTypes.ROW(
                        DataTypes.FIELD("tran_date", DataTypes.STRING()),
                        DataTypes.FIELD("entity_id", DataTypes.BIGINT()),
                        DataTypes.FIELD("amt", DataTypes.STRING())
                ))
                .build())
        .format("protobuf")
        .option("protobuf.message-class-name", RawDataCarrierProto.class.getName())
        .option("topic", topicName)
        .option("properties.bootstrap.servers", kafkaBootstrapServers)
        .option("properties.group.id", consumerGroupId)
        .option("scan.startup.mode", "latest-offset")
        .build();

解决方案

原因分析

问题核心在于Flink生成水位线代码时,无法正确处理java.time.LocalDateTime类型的UDF返回值——水位线机制依赖Flink原生时间类型(如TIMESTAMP_LTZ、TIMESTAMP),自定义UDF返回LocalDateTime会导致Janino编译阶段无法解析类型依赖;同时在columnByExpression中嵌套UDF与CAST操作,也会干扰Flink对事件时间字段的合法性识别。

修复方案

方案1:修改UDF返回Flink兼容的时间类型

将UDF返回类型改为java.sql.Timestamp,直接输出Flink可识别的时间类型,避免中间转换:

public class ToLocalDateTimeFunction extends ScalarFunction {
    private static final DateTimeFormatter format = new DateTimeFormatterBuilder()
            .appendPattern("yyyyMMddHHmmss")
            .appendValue(ChronoField.MILLI_OF_SECOND, 3)
            .toFormatter();

    public Timestamp eval(String value) {
        try {
            LocalDateTime localDateTime = LocalDateTime.parse(value, format);
            return Timestamp.valueOf(localDateTime);
        } catch (Exception e) {
            return null;
        }
    }
}

之后在TableDescriptor中无需额外CAST,直接使用UDF结果配置水位线:

.columnByExpression("rowtime", "TO_LOCALDATETIME(raw_data.tran_date)")
.watermark("rowtime", "rowtime - INTERVAL '10' MINUTE")

方案2:使用Flink内置函数替代自定义UDF

Flink 1.16+版本的TO_TIMESTAMP已支持yyyyMMddHHmmssSSS格式解析,直接使用原生函数即可:

.columnByExpression("rowtime", "CAST(TO_TIMESTAMP(raw_data.tran_date, 'yyyyMMddHHmmssSSS') AS TIMESTAMP_LTZ(3))")
.watermark("rowtime", "rowtime - INTERVAL '10' MINUTE")

此方案可完全规避自定义UDF带来的编译问题。

方案3:调整Schema字段声明顺序

将raw_data字段的声明移至最前面,确保Flink解析表达式时能正确识别原始字段类型:

.schema(Schema.newBuilder()
        .column("raw_data", DataTypes.ROW(
                DataTypes.FIELD("tran_date", DataTypes.STRING()),
                DataTypes.FIELD("entity_id", DataTypes.BIGINT()),
                DataTypes.FIELD("amt", DataTypes.STRING())
        ))
        .columnByExpression("entityId", "raw_data.entity_id")
        .columnByExpression("tranDate", "TO_LOCALDATETIME(raw_data.tran_date)")
        .columnByExpression("rowtime", "CAST(TO_LOCALDATETIME(raw_data.tran_date) AS TIMESTAMP_LTZ(3))")
        .columnByExpression("amt", "TO_BIGDECIMAL(raw_data.amt)")
        .watermark("rowtime", "rowtime - INTERVAL '10' MINUTE")
        .columnByExpression("processTime", "PROCTIME()")
        .build())

内容的提问来源于stack exchange,提问作者Süleyman Fazıl Yeşil

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 18:24:52