乐于分享
好东西不私藏

Flink源码学习系列课程-05-flatMap方法与OneInputTransformation的诞生

Flink源码学习系列课程-05-flatMap方法与OneInputTransformation的诞生

核心文件DataStream.java

从哪里调用

你在 WordCount 里调用了 text.flatMap(new Tokenizer())text 是上一课讲的 DataStream,它的 transformation 指向 SourceTransformation(id=1)。flatMap 的内部发生了什么?

这一课解决什么问题

这一行代码:

text.flatMap(new Tokenizer()).name("tokenizer")

从你写的 Tokenizer 对象到 Flink 内部的 OneInputTransformation,中间经历了 5 层包装。不是 Flink 喜欢过度工程化——每一层都有明确的设计理由。我们逐层剥开。

第一关:flatMap 入口——类型推断的魔法

打开 DataStream.java,跳到第 458-464 行:

public <R> SingleOutputStreamOperator<R> flatMap(FlatMapFunction<T, R> flatMapper){// ① 类型推断:从 FlatMapFunction 的泛型参数中提取输出类型    TypeInformation<R> outType =            TypeExtractor.getFlatMapReturnTypes(                    clean(flatMapper), getType(), Utils.getCallLocationName(), true);// ② 调用重载版本return flatMap(flatMapper, outType);}

TypeExtractor.getFlatMapReturnTypes() 通过 Java 反射,从 FlatMapFunction<String, Tuple2<String, Integer>> 的泛型签名中提取 Tuple2<String, Integer> 的类型信息。

为什么需要这一层? Flink 需要知道流中每个元素的类型,因为它要用 TypeSerializer 做序列化——把 Java 对象变成网络上的字节流。没有类型信息,Serializer 不知道该序列化什么样的对象。如果反射失败(比如用了 Lambda),方法会抛出 InvalidTypesException,提醒你显式指定类型。

Utils.getCallLocationName() 返回调用 flatMap() 的代码位置(类名:行号),用于错误日志中定位问题。

第二关:StreamFlatMap——为什么不能直接跑你的 Tokenizer

跳到第 479-482 行:

public <R> SingleOutputStreamOperator<R> flatMap(        FlatMapFunction<T, R> flatMapper, TypeInformation<R> outputType){return transform("Flat Map", outputType, new StreamFlatMap<>(clean(flatMapper)));}

StreamFlatMap 是 Flink 对 FlatMapFunction 的运行时包装:

StreamFlatMap extends AbstractStreamOperator<OUT>              implements OneInputStreamOperator<IN, OUT>

它的核心方法(简化版):

// StreamFlatMap.javapublicvoidprocessElement(StreamRecord<IN> element)throws Exception {    userFunction.flatMap(element.getValue(), output);  // 直接调你的 Tokenizer}

关键细节①:Flink 为什么不直接跑你写的 Tokenizer 而要包装一层?因为运行时不只需要"处理逻辑",还需要:

  • Metrics 收集(这个算子每秒处理了多少条数据、延迟多少)
  • Checkpoint 协调(收到 barrier 时暂停处理)
  • Watermark 处理(事件时间推进)
  • 生命周期管理(算子 open/close 时初始化/释放资源)

如果让你的 Tokenizer 直接实现所有这些接口,用户代码会膨胀到无法维护。包装一层 = 运行时胶水与业务逻辑分离。

第三关:transform()——SimpleOperatorFactory 的角色

跳到第 793-799 行:

public <R> SingleOutputStreamOperator<R> transform(        String operatorName,        TypeInformation<R> outTypeInfo,        OneInputStreamOperator<T, R> operator){return doTransform(operatorName, outTypeInfo, SimpleOperatorFactory.of(operator));}

SimpleOperatorFactory.of(operator) 创建一个 StreamOperatorFactory

为什么这里用工厂而不用直接 new? 你是在 Client 端调用 flatMap——这里没有 StateBackend、没有网络环境、没有 Checkpoint Storage。真正的算子实例在 TM 端创建才能拿到这些运行时依赖。工厂可以序列化后通过网络发送到 TM,在 TM 侧调用 factory.createStreamOperator() 创建算子。这就是工厂模式在分布式系统中的意义。

第四关(关键!):doTransform——惰性求值的终极证明

跳到第 823-847 行,这是整个 flatMap 调用链的心脏:

protected <R> SingleOutputStreamOperator<R> doTransform(        String operatorName,        TypeInformation<R> outTypeInfo,        StreamOperatorFactory<R> operatorFactory){// ① 触发输入类型校验:如果类型缺失,在这里报错(Fast Fail)    transformation.getOutputType();// ② ★ 创建新的 Transformation 节点——连接父节点和当前节点    OneInputTransformation<T, R> resultTransform =new OneInputTransformation<>(this.transformation,      // ← 父 Transformation(SourceTransformation id=1)                    operatorName,             // "Flat Map"                    operatorFactory,          // SimpleOperatorFactory(StreamFlatMap(Tokenizer))                    outTypeInfo,              // TypeInformation<Tuple2<String,Integer>>                    environment.getParallelism(),  // 继承环境的默认并行度false);// ③ 创建新的 DataStream 壳返回给用户    SingleOutputStreamOperator<R> returnStream =new SingleOutputStreamOperator(environment, resultTransform);// ④ ★ 把新 Transformation 注册到环境的 transformations List    getExecutionEnvironment().addOperator(resultTransform);return returnStream;}

关键细节②:并行度继承

environment.getParallelism()  // 使用环境的默认并行度

如果你之前调用了 env.setParallelism(4),所有新算子都继承 4。但如果你在链式调用中写 .flatMap().setParallelism(2),这个值在 doTransform() 之后通过第③步拿到 SingleOutputStreamOperator,再通过 .setParallelism(2) 覆盖——还记得上一课的配置优先级吗?后写覆盖先写。

关键细节③:addOperator——惰性求值的终极证明

Ctrl/Cmd + 点击addOperator,你会发现:

// StreamExecutionEnvironmentpublicvoidaddOperator(Transformation<?> transformation){    Preconditions.checkNotNull(transformation, "transformation must not be null.");this.transformations.add(transformation);}

就一行 ArrayList.add()! 没有启动线程,没有创建网络连接,没有分配内存。只是把新节点加到 List 里。这是惰性求值的终极证明——你所有算子调用只是在 ArrayList 里追加元素,直到你调用 env.execute() 才真正开始执行。

实践环节

在 addOperator 的 transformations.add(transformation) 处打断点。运行 WordCount,你会发现每调一次 flatMap/keyBy/sum,断点触发一次。再次确认:所有算子调用只是往 ArrayList 追加元素。

第五关:ClosureCleaner——Lambda 的序列化保险

clean(flatMapper)  // 在 flatMap() 入口就已经调用了

对于 Lambda 表达式 value -> value.toLowerCase().split("\\W+")

  • 如果 Lambda 捕获了外部变量(如 separator 字段),ClosureCleaner 检查该变量是否 Serializable
  • 如果不是 Serializable,提交作业时报 NotSerializableException(Fast Fail)
  • 如果可以序列化,处理为跨 JVM 兼容的版本

为什么 Lambda 需要这个处理? Java 的 Lambda 默认序列化行为取决于 JVM 实现——同一个 Lambda 在 HotSpot 和 OpenJ9 上序列化结果可能不同。Flink 的作业要在不同 JVM 之间传递,必须保证序列化兼容性。

调用链追踪

text.flatMap(new Tokenizer())                              // DataStream.java:458  │  ├─→ clean(flatMapper)                                   // ClosureCleaner  │  ├─→ TypeExtractor.getFlatMapReturnTypes(...)            // 类型推断  │       └→ 反射读取 FlatMapFunction<String, Tuple2<...>> 的泛型参数  │  ├─→ new StreamFlatMap<>(clean(flatMapper))              // DataStream.java:481  │       └→ 内部包装,添加 processElement() 的胶水代码  │  ├─→ transform(name, type, streamFlatMap)                // DataStream.java:793  │       └→ SimpleOperatorFactory.of(streamFlatMap)      // 工厂化  │  └─→ doTransform(name, type, factory)                    // DataStream.java:823        │        ├─→ transformation.getOutputType()                // 类型校验(Fast Fail)        │        ├─→ new OneInputTransformation<>(                 // 创建 DAG 节点        │       this.transformation,    // parent = SourceTransformation(id=1)        │       "Flat Map",             // name        │       factory,                // SimpleOperatorFactory        │       outTypeInfo,            // Tuple2<String,Integer>        │       env.getParallelism(),   // 默认并行度        │       false)        │        ├─→ new SingleOutputStreamOperator(env, result)   // 创建返回对象        │        └─→ env.addOperator(resultTransform)              // ★ ArrayList.add()              └→ transformations.add(transform)

总结

┌──────────────────────────────────────────────────────────────────┐│  flatMap() 的 5 层包装:                                          ││                                                                   ││  1. ClosureCleaner.clean()        ← 让函数可跨 JVM 序列化         ││  2. StreamFlatMap                 ← 添加运行时胶水代码             ││  3. SimpleOperatorFactory         ← 工厂化(延迟创建到 TM 侧)     ││  4. OneInputTransformation        ← 创建 DAG 节点(input 连接上游)││  5. env.addOperator()             ← ArrayList.add() — 惰性求值     ││                                                                   ││  核心设计模式:                                                    ││  • 装饰器模式:StreamFlatMap 装饰用户的 FlatMapFunction             ││  • 工厂模式:SimpleOperatorFactory 延迟创建到 TM 侧                ││  • 单向链表:每个 Transformation.input 指向上游                    │└──────────────────────────────────────────────────────────────────┘

下一课预告

keyBy() 为什么不一样?它不创建新的 Transformation,只返回 KeyedStream。分区效果在 sum() 时才体现——ReduceTransformation 内部自动创建 PartitionTransformation