使用自定义UDF转换字段配置Flink水位线时程序编译失败
因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

