Flink 数据类型与序列化从类型体系到高性能序列化
1. 为什么 Flink 要“管类型”?
Flink 在计划与优化阶段需要精确的类型信息(TypeInformation),以便选择更高效的内存布局、序列化器与算子实现。了解得越多,运行时就能少做“猜测”,序列化更快、内存更省、排序/KeyBy 等操作也能更高效。
2. Flink 支持的数据类型(8 大类)
- 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 下标一致)。
- 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()进行单测校验。
- Primitive Types
- Java 基本类型及其包装类、
String、Double等,原生支持。
- 常见集合类型(Map/List/Set/Collection)
高效专用序列化,但有两条硬性要求:
- 具体泛型:
List<String>✅,List<?>/List<T>/List❌; - 接口类型:
List<String>✅,LinkedList<String>❌(运行中不保证实现类)。
否则按“通用类”处理,必要时注册自定义序列化器。
- General Class Types(通用类)
- 不能被识别为 POJO 的 Java 类,Flink 视为黑箱,用 Kryo 序列化。
- 不适合包含文件句柄、IO 流、Native 资源等不可序列化字段。
- Values(自定义序列化)
- 实现
org.apache.flink.types.Value,手写read/write,可极致优化(如稀疏向量仅写非零项)。 CopyableValue可自定义 copy 逻辑。- Flink 内置
IntValue/StringValue/...等“可变值类型”,降低 GC 压力。
- Hadoop Writables
- 兼容
org.apache.hadoop.Writable,直接用其write/readFields。
- 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. 最常见问题与解决方案(速查)
- 集合类型不生效 / 序列化退化
- 使用 接口 + 具体泛型:
List<String>✅,LinkedList<String>/List<?>❌。
- 类型推断失败 / 运行时报 Kryo
-
给函数链路添加
returns(...)Type Hint; -
或实现 ResultTypeQueryable;
-
或禁用 Kryo 回退,逼出问题类型并补齐自定义序列化器:
pipeline.generic-types: false
- Kryo 无法处理三方类型
- 在
pipeline.serialization-config注册合适的 Kryo 序列化器; - 或改造为 POJO/Value 类型。
- 想彻底避开 Kryo
- 尽量让类型成为 POJO/基本类型/支持的集合;
- 必要时使用 Value 自实现
read/write; - 通过
pipeline.generic-types: false保持“强约束”。
- 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快速定位不透明类型。
更多推荐


所有评论(0)