【问题标题】:How to efficiently serialize POJO with LocalDate field in Flink?如何在 Flink 中使用 LocalDate 字段有效地序列化 POJO?
【发布时间】:2021-02-17 08:52:06
【问题描述】:

我们的一些 POJO 包含来自 java.time API 的字段(LocalDate、LocalDateTime)。当我们的管道处理它们时,我们可以在日志中看到以下信息:

org.apache.flink.api.java.typeutils.TypeExtractor - Class class java.time.LocalDate cannot be used as a POJO type because not all fields are valid POJO fields, and must be processed as GenericType. Please read the Flink documentation on "Data Types & Serialization" for details of the effect on performance.

据我了解,LocalDate 不能归类为 POJO,因此 flink 不使用 POJO 序列化程序,而是退回到效率较低的 Kryo。然而,由于 1.9.0 版本的 flink 为 java.time 类提供了专用的序列化器(例如LocalDateSerializer),所以我希望这些序列化器可以在这里完成工作,从而允许 POJO 序列化器用于我们的类。不是这样吗?如果是,是否有任何性能影响?如果不是,这种情况的最佳解决方案是什么?

在项目中,我们使用 Flink 1.11 和 Java 1.8。

【问题讨论】:

    标签: java apache-flink


    【解决方案1】:

    由于向后兼容,即使在 Flink 中引入了新的序列化器,它也不能自动使用。但是,您可以告诉 Flink 像这样将它用于您的 POJO(如果您开始时没有使用 Kryo 的先前保存点):

    @TypeInfo(MyClassTypeInfoFactory.class)
    public class MyClass {
      public int id;
      public LocalDate date;
      // ...
    }
    
    public class MyClassTypeInfoFactory extends TypeInfoFactory<MyClass> {
      @Override
      public TypeInformation<MyClass> createTypeInfo(
          Type t, Map<String, TypeInformation<?>> genericParameters) {
    
        return Types.POJO(MyClass.class, new HashMap<String, TypeInformation<?>>() { {
            put("id", Types.INT);
            put("date", Types.LOCAL_DATE);
            // ...
        } } );
      }
    }
    

    您必须为您的所有 POJO 字段提供类型,如图所示,但您可以使用 Types 类中的许多帮助程序。此外,通过像这样使用TypeInfoFactory,您不必担心 Flink 使用这种类型的所有地方 - 它总是会派生给定的类型信息。

    如果您需要转换旧的保存点以使用新的序列化程序,您可能还需要查看Flink's State Processor API

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-09-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多