Flink环境下如何确认自定义Java类被识别为POJO类还是通用类
确认Flink对类的类型识别结果的方法
- 直接通过TypeInformation判断
执行如下代码即可直接拿到识别结果:TypeInformation<Speed> typeInfo = TypeInformation.of(Speed.class); // 输出是否为POJO类型 System.out.println("是否为POJO:" + (typeInfo instanceof PojoTypeInfo)); // 输出实际的类型分类 System.out.println("类型分类:" + typeInfo.getClass().getSimpleName()); // 输出对应使用的序列化器 System.out.println("序列化器实现:" + typeInfo.createSerializer(new ExecutionConfig())); - 查看类型分析日志
将org.apache.flink.api.java.typeutils包的日志级别调整为DEBUG后提交作业,Flink会输出完整的类型校验过程,不符合POJO规则的原因也会明确打印。 - 查看作业Web UI
作业运行后在对应算子的详情页,可直接查看输入/输出的类型标注,明确显示是POJO还是通用类型。
你当前场景的识别结果说明
Flink判定POJO的规则如下:
- 类为公共非抽象类
- 存在公共无参构造方法
- 所有非静态、非transient字段要么是公共的,要么符合JavaBean的getter/setter规范
- 所有字段的类型都有对应的序列化实现
你的Speed类本身满足前3个规则,所以会被识别为POJO类,仅内部的telemetry字段因为Telemetry类型不符合POJO规则,会作为通用类型用Kryo序列化。
另外你当前的注册代码存在错误:你自定义的是Telemetry类型的序列化器,但注册代码绑定的是DataElement类和对应的序列化器,没有注册Telemetry类,所以你查询已注册Kryo类型会返回
env.getConfig().registerTypeWithKryoSerializer(Telemetry.class, TelemetryKryoSerializer.class);
修正注册后你写的自定义序列化器就会生效。
内容的提问来源于stack exchange,提问作者Ahmed A
相关产品推荐
相关产品推荐

