Flink 1.17.2中POJO测试通过但执行报错问题排查求助
问题定位:Flink 1.17.2中Lombok注解POJO测试通过但流处理任务报错
问题重现
定义的Lombok注解POJO类
@Data @AllArgsConstructor @NoArgsConstructor public class IdCount { private Integer id; private String name; }
POJO验证测试代码
public class Test { public static boolean isPojoClass(Class<?> pojoClass) { try { TypeInformation<?> typeInfo = TypeExtractor.createTypeInfo(pojoClass); return typeInfo instanceof PojoTypeInfo; } catch (Exception e) { return false; } } public static void main(String[] args) { boolean isPojo = isPojoClass(IdCount.class); if (isPojo) { System.out.println("This is a POJO class."); } else { System.out.println("This is not a POJO class."); } } }
测试输出
This is a POJO class.
实际流处理任务代码
public class TimeWindowExample { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env); env.setParallelism(1); tableEnv.executeSql("CREATE TABLE dataGen (\n" + " id INT,\n" + " name STRING\n" + ") WITH (\n" + " 'connector' = 'datagen',\n" + " 'rows-per-second'='1',\n" + " 'fields.id.kind'='random',\n" + " 'fields.id.min'='1',\n" + " 'fields.id.max'='10',\n" + " 'fields.name.length'='10'\n" + ")"); Table table = tableEnv.sqlQuery("select * from dataGen"); DataStream<IdCount> dataStream = tableEnv.toDataStream(table, IdCount.class); dataStream .countWindowAll(3) .max("id") .print() ; env.execute(); } }
报错信息
Exception in thread "main" org.apache.flink.api.common.typeutils.CompositeType$InvalidFieldReferenceException: Cannot reference field by field expression on *cn.chatdoge.flink117.POJO.IdCount<`id` INT, `name` STRING>* (cn.chatdoge.flink117.POJO.IdCount, org.apache.flink.table.runtime.typeutils.ExternalSerializer) Field expressions are only supported on POJO types, tuples, and case classes. (See the Flink documentation on what is considered a POJO.) at org.apache.flink.streaming.util.typeutils.FieldAccessorFactory.getAccessor(FieldAccessorFactory.java:256) at org.apache.flink.streaming.api.functions.aggregation.ComparableAggregator.<init>(ComparableAggregator.java:78) at org.apache.flink.streaming.api.datastream.AllWindowedStream.max(AllWindowedStream.java:1366) at cn.chatdoge.flink117.window.TimeWindowExample.main(TimeWindowExample.java:55)
原因分析
测试代码与实际任务的POJO识别逻辑存在差异:
- 测试代码直接通过
TypeExtractor.createTypeInfo(IdCount.class)生成PojoTypeInfo,这是DataStream API原生的POJO识别逻辑,Lombok生成的无参构造、getter/setter完全符合Flink POJO的要求,因此测试通过。 - 但通过
tableEnv.toDataStream(table, IdCount.class)转换时,Table API使用ExternalSerializer而非Flink原生POJO序列化器处理IdCount。此时DataStream的TypeInformation并非PojoTypeInfo,导致后续窗口聚合的max("id")无法识别字段表达式,触发报错。
解决方案
方案1:显式指定POJO类型信息
在转换时手动生成并指定PojoTypeInfo,确保DataStream识别为POJO类型:
// 方式1:手动转换并指定TypeInformation DataStream<IdCount> dataStream = tableEnv.toDataStream(table) .map(row -> new IdCount(row.getFieldAs("id"), row.getFieldAs("name"))) .returns(TypeInformation.of(IdCount.class)); // 方式2:直接传入PojoTypeInfo PojoTypeInfo<IdCount> pojoTypeInfo = (PojoTypeInfo<IdCount>) TypeExtractor.createTypeInfo(IdCount.class); DataStream<IdCount> dataStream = tableEnv.toDataStream(table, pojoTypeInfo);
方案2:手动编写POJO的getter/setter(可选)
若Lombok生成的字节码存在识别异常,可手动编写符合Flink要求的POJO代码,避免依赖Lombok的自动生成:
public class IdCount { private Integer id; private String name; public IdCount() {} public IdCount(Integer id, String name) { this.id = id; this.name = name; } public Integer getId() { return id; } public void setId(Integer id) { this.id = id; } public String getName() { return name; } public void setName(String name) { this.name = name; } }
方案3:在Table API阶段完成聚合
直接在Table层面执行聚合操作,再转换为DataStream,规避DataStream侧的POJO识别问题:
Table resultTable = tableEnv.sqlQuery("SELECT MAX(id) FROM dataGen GROUP BY TUMBLE(ROWS, INTERVAL '3' ROWS)"); tableEnv.toDataStream(resultTable).print();
内容的提问来源于stack exchange,提问作者Samuel Mau
相关产品推荐
相关产品推荐

