1. 为什么 Flink 要“管类型”?

Flink 在计划与优化阶段需要精确的类型信息(TypeInformation),以便选择更高效的内存布局、序列化器与算子实现。了解得越多,运行时就能少做“猜测”,序列化更快、内存更省、排序/KeyBy 等操作也能更高效。

2. Flink 支持的数据类型(8 大类)

  1. Java Tuples
  • Tuple1 ~ Tuple25 固定长度、可嵌套:

    DataStream<Tuple2<String, Integer>> s = env.fromElements(
        new Tuple2<>("hello", 1), new Tuple2<>("world", 2));
    s.map((MapFunction<Tuple2<String,Integer>, Integer>) v -> v.f1);
    s.keyBy(v -> v.f0);
    
  • 字段从 f0 开始(与 Java 下标一致)。

  1. Java POJOs(重点)
    满足以下条件即被识别为 POJO(从而“按字段”处理,性能更好):
  • public 类,无参 public 构造
  • 字段 public 或有 public getter/setter(Java Beans 规范);
  • 字段类型有可用的序列化器。
    示例:
public class WordWithCount {
  public String word;
  public int count;
  public WordWithCount() {}
  public WordWithCount(String word, int count) { this.word = word; this.count = count; }
}
DataStream<WordWithCount> s = env.fromElements(
  new WordWithCount("hello", 1), new WordWithCount("world", 2));
s.keyBy(v -> v.word);
  • 默认由 PojoSerializer 处理(必要时 Kryo 兜底)。若是 Avro Specific/Reflect 类型,则用 AvroSerializer
  • 可用 PojoTestUtils.assertSerializedAsPojo() 进行单测校验。
  1. Primitive Types
  • Java 基本类型及其包装类、StringDouble 等,原生支持。
  1. 常见集合类型(Map/List/Set/Collection)
    高效专用序列化,但有两条硬性要求
  • 具体泛型List<String> ✅,List<?>/List<T>/List ❌;
  • 接口类型List<String> ✅,LinkedList<String> ❌(运行中不保证实现类)。
    否则按“通用类”处理,必要时注册自定义序列化器。
  1. General Class Types(通用类)
  • 不能被识别为 POJO 的 Java 类,Flink 视为黑箱,用 Kryo 序列化。
  • 不适合包含文件句柄、IO 流、Native 资源等不可序列化字段。
  1. Values(自定义序列化)
  • 实现 org.apache.flink.types.Value,手写 read/write,可极致优化(如稀疏向量仅写非零项)。
  • CopyableValue 可自定义 copy 逻辑。
  • Flink 内置 IntValue/StringValue/... 等“可变值类型”,降低 GC 压力。
  1. Hadoop Writables
  • 兼容 org.apache.hadoop.Writable,直接用其 write/readFields
  1. Special Types
  • Java API 的 Either<L,R>(类似 Scala),常用于分叉输出或错误处理。

3. Java 的类型擦除与 Flink 的类型推断

  • Java 在编译后擦除泛型DataStream<String>DataStream<Long> 在 JVM 层面看起来一样。

  • Flink 在提交前(main 执行期间)尽力反射还原类型信息,并存入 TypeInformation

    DataStream<?> ds = ...;
    TypeInformation<?> t = ds.getType(); // Flink 内部类型描述
    
  • 推断有边界:如 fromCollection() 或泛型函数 MapFunction<I,O> 复杂场景,需要开发者配合(传入 Type Hint)。

常见“协助”手段

  • ResultTypeQueryable:格式/函数实现该接口,显式返回结果类型。
  • Type Hints:为 Java API 提供类型提示(下一节详述)。

4. Type Hints 与 TypeInformation/Serializer 创建

4.1 Type Hints(Java API)

当推断失败时,显式声明返回类型:

DataStream<SomeType> result = stream
  .map(new MyGenericFunction<Long, SomeType>())
  .returns(SomeType.class);

// 对于泛型复合类型,用 TypeHint:
DataStream<Tuple2<Integer, SomeType>> result2 = stream
  .map(...)
  .returns(new TypeHint<Tuple2<Integer, SomeType>>(){});

4.2 创建 TypeInformation

// 非泛型
TypeInformation<String> info1 = TypeInformation.of(String.class);
// 泛型(捕获匿名子类上的泛型参数)
TypeInformation<Tuple2<String, Double>> info2 =
    TypeInformation.of(new TypeHint<Tuple2<String, Double>>(){});

4.3 获取 TypeSerializer 的两种方式

// 方式一:通过 TypeInformation
TypeSerializer<T> ser = typeInfo.createSerializer(serializerConfig);

// 方式二:RichFunction 内,从 RuntimeContext 获取
TypeSerializer<T> ser2 = getRuntimeContext().createSerializer(typeInfo);

4.4 Java 8 Lambda 的类型提取

  • Flink 尝试从 lambda 目标方法的泛型签名获取类型;
  • 若编译器未生成签名或推断异常,请用 returns(...) 明确指定

5. POJO 序列化的控制:Avro/Kryo 强制策略

  • 默认:POJO 用 PojoSerializer,未知类型回退 Kryo。

  • 强制 Avro(包含 flink-avro 模块):

    pipeline.force-avro: true
    
  • 强制 Kryo

    pipeline.force-kryo: true
    
  • 禁用 Kryo 回退(希望确保所有类型都有“高效/可控”的序列化器):

    pipeline.generic-types: false
    

    遇到需要 Kryo 的类型将直接抛错,方便你补齐类型注册与自定义序列化器。

6. 注册子类型 & 自定义序列化器

6.1 注册子类型/序列化(YAML)

pipeline.serialization-config:
  # 为 POJO 注册自定义序列化器
  - org.example.MyCustomType1: {type: pojo, class: org.example.MyCustomSerializer1}
  # 为通用类型注册 Kryo 序列化器
  - org.example.MyCustomType2: {type: kryo, kryo-type: registered, class: org.example.MyCustomSerializer2}

6.2 代码中设置

Configuration config = new Configuration();
config.set(PipelineOptions.SERIALIZATION_CONFIG, List.of(
  "org.example.MyCustomType1: {type: pojo, class: org.example.MyCustomSerializer1}",
  "org.example.MyCustomType2: {type: kryo, kryo-type: registered, class: org.example.MyCustomSerializer2}"
));
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(config);

说明:大量三方库(如 Guava 的集合)默认 Kryo 不总是“开箱即用”,建议为问题类型注册合适的 Kryo 序列化器。

7. 用 TypeInfoFactory 自定义类型信息(高级)

当 Flink 的默认类型提取不满足需求,可用 TypeInfoFactory 插桩,向 Flink 类型系统提供自定义 TypeInformation

  • 通过 配置 关联:

    pipeline.serialization-config:
      - org.example.MyCustomType: {type: typeinfo, class: org.example.MyCustomTypeInfoFactory}
    
  • 注解 关联:

    @TypeInfo(MyTupleTypeInfoFactory.class) // 标注在类型或 POJO 字段上
    public class MyTuple<T0, T1> { public T0 myfield0; public T1 myfield1; }
    

工厂示例:

public class MyTupleTypeInfoFactory extends TypeInfoFactory<MyTuple> {
  @Override
  public TypeInformation<MyTuple> createTypeInfo(
      Type t, Map<String, TypeInformation<?>> genericParameters) {
    return new MyTupleTypeInfo(genericParameters.get("T0"), genericParameters.get("T1"));
  }
}

若类型含泛型参数且需与输入类型互推,请实现 TypeInformation#getGenericParameters,实现双向映射

8. 最常见问题与解决方案(速查)

  1. 集合类型不生效 / 序列化退化
  • 使用 接口 + 具体泛型List<String> ✅,LinkedList<String>/List<?> ❌。
  1. 类型推断失败 / 运行时报 Kryo
  • 给函数链路添加 returns(...) Type Hint

  • 或实现 ResultTypeQueryable

  • 禁用 Kryo 回退,逼出问题类型并补齐自定义序列化器:

    pipeline.generic-types: false
    
  1. Kryo 无法处理三方类型
  • pipeline.serialization-config 注册合适的 Kryo 序列化器;
  • 或改造为 POJO/Value 类型。
  1. 想彻底避开 Kryo
  • 尽量让类型成为 POJO/基本类型/支持的集合
  • 必要时使用 Value 自实现 read/write
  • 通过 pipeline.generic-types: false 保持“强约束”。
  1. POJO 不被识别
  • 检查:public 类、无参构造、字段可访问或 bean 风格 getter/setter、非 static/非 transient。
  • 否则会退化为 GenericType → Kryo。

9. 实用代码/配置片段(可直接用)

9.1 快速判断某类是否按 POJO 序列化(单测)

import static org.apache.flink.types.PojoTestUtils.assertSerializedAsPojo;
public class PojoTest {
  @Test
  public void testPojo() {
    assertSerializedAsPojo(WordWithCount.class);
  }
}

9.2 为通用类型注册 Kryo 序列化器(YAML)

pipeline.serialization-config:
  - com.google.common.collect.ImmutableList: {type: kryo, kryo-type: registered, class: com.esotericsoftware.kryo.serializers.CollectionSerializer}

9.3 强制 Avro 处理 POJO(全局开关)

pipeline.force-avro: true
# 加入 flink-avro 依赖

9.4 在 Map 中添加 Type Hint(泛型返回)

DataStream<Tuple2<String, Long>> out = in
  .map(new AppendOne<>())
  .returns(new TypeHint<Tuple2<String, Long>>(){});

10. 性能与工程化建议

  • 优先 POJO:字段可见、稳定 schema、便于优化与演进;
  • 集合用接口 + 具体泛型:让 Flink 用上专用序列化;
  • 减少对象创建:可用 Value 可变类型降 GC;
  • 对热点类型做序列化基准:从 Kryo → Avro/Value/自定义序列化迁移;
  • 在关键路径显式声明 Type Hint:避免“隐式推断”成为性能黑盒;
  • 有问题就关掉 Kryo 回退pipeline.generic-types: false 快速定位不透明类型。
Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐