核心文件: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。
夜雨聆风