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

Flink 1.17.2中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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 11:55:03