核心文件:DataStream.java, SingleOutputStreamOperator.java, KeyedStream.java
从哪里调用
上一课 execute() 讲解了惰性求值结束后的执行过程。现在我们要倒回去——在 execute() 之前,你的 flatMap()、keyBy()、sum() 调用到底做了什么?
这一课解决什么问题
你的 WordCount 代码里,不同方法返回了不同的类型:
DataStream<String> text = env.fromSource(...); // DataStreamSingleOutputStreamOperator words = text.flatMap(new Tokenizer()); // SingleOutputStreamOperatorKeyedStream<Tuple2<String, Integer>, String> keyed = ...keyBy(...); // KeyedStreamSingleOutputStreamOperator counts = keyed.sum(1); // 又回到 SingleOutputStreamOperator你注意到没有——每次调完,返回的类型都不一样。这背后是 Flink 的设计哲学:不同返回类型 = 编译期的可用方法集控制。flatMap 之后你可以设置并行度,keyBy 之后你才能调 sum/reduce/window。如果 keyBy 返回普通的 DataStream,你在编译期就会犯"在非 Keyed 流上调用 sum"的错误。
第一关:DataStream 的核心——两个 final 字段
打开 DataStream.java,跳到第 111-129 行:
@PublicpublicclassDataStream<T> {protectedfinal StreamExecutionEnvironment environment; // 环境引用protectedfinal Transformation<T> transformation; // ★ DAG 节点的句柄publicDataStream(StreamExecutionEnvironment environment, Transformation<T> transformation){this.environment = Preconditions.checkNotNull(environment);this.transformation = Preconditions.checkNotNull(transformation); }}就两个字段。就这么简单。
DataStream = 壳,Transformation = 核心。 每个 DataStream 对象只是对环境引用和 Transformation 引用的封装。你调用的 flatMap()、keyBy() 等方法,本质上都是在操作或创建新的 Transformation。
为什么是 final? DataStream 一旦创建,其对应的 Transformation 就不能改变了。DataStream<String> text 始终代表文件源,不管你在后面调了多少次 .flatMap(),text 本身不变——.flatMap() 返回的是一个新的 DataStream 对象。这是一种不可变设计,防止你在链式调用中产生副作用。
实践环节
在 IDE 里,Ctrl/Cmd + 点击transformation 字段,看它是什么类型。你会发现它可以是 SourceTransformation、OneInputTransformation、ReduceTransformation 等多种子类——这就是多态的核心。
第二关:getId()——DAG 节点的唯一身份证
DataStream.java 第 137-139 行:
@InternalpublicintgetId(){return transformation.getId();}每次 flatMap() 创建新 Transformation 时,ID 由 StreamExecutionEnvironment 内的计数器自增生成。这个 ID 是整个 DAG 中节点的唯一标识,在后续的 StreamGraph、JobGraph 中一直沿用。
为什么需要 ID? 你在大图的上下文中,有 50 个节点。你需要一种方式来标识"从 Tokenizer(id=2) 到 Counter(id=4) 有一条边"。ID 就是这个引用的锚点。
第三关(关键!):SingleOutputStreamOperator——壳的升级版
打开 SingleOutputStreamOperator.java,跳到第 49-64 行:
@PublicpublicclassSingleOutputStreamOperator<T> extendsDataStream<T> {protectedboolean nonParallel = false; // 强制单并行度的标记private Map<OutputTag<?>, TypeInformation<?>> requestedSideOutputs = new HashMap<>();protectedSingleOutputStreamOperator( StreamExecutionEnvironment environment, Transformation<T> transformation){super(environment, transformation); // 直接传给 DataStream 的构造函数 }}继承 DataStream,没有新增核心字段——nonParallel 和 requestedSideOutputs 都是辅助性的。它的价值在于方法:
.name("tokenizer") | transformation.setName() | |
.setParallelism(4) | transformation.setParallelism() | |
.uid("my-op") | transformation.setUid() | |
.disableChaining() | transformation.setChainingStrategy(NEVER) | |
.startNewChain() | transformation.setChainingStrategy(HEAD) | |
.slotSharingGroup("g") | transformation.setSlotSharingGroup() | |
.setBufferTimeout(100) | transformation.setBufferTimeout() |
所有这些方法返回 this(SingleOutputStreamOperator<T>),形成 Fluent API:
text.flatMap(new Tokenizer()) .name("tokenizer") .setParallelism(4) .uid("tokenizer-v2");为什么 flatMap 返回它而不是普通的 DataStream? 因为 flatMap 创建了新算子——你接下来要设置算子的名称、并行度、Chain 策略。SingleOutputStreamOperator 正是提供了这些配置方法。如果返回普通的 DataStream,你就只能在每次 flatMap 之后手动去 env 上翻找最新的 Transformation——那就太反人类了。
第四关:disableChaining 和 startNewChain——Facade 模式的优雅
SingleOutputStreamOperator.java 第 275-289 行:
@PublicEvolvingpublic SingleOutputStreamOperator<T> disableChaining(){return setChainingStrategy(ChainingStrategy.NEVER); // ← 前后都不 Chain}@PublicEvolvingpublic SingleOutputStreamOperator<T> startNewChain(){return setChainingStrategy(ChainingStrategy.HEAD); // ← 不和前驱 Chain,但可以和后继 Chain}内部的 setChainingStrategy() 是 private 的——你不能直接设置任意策略。
为什么? 这是 Facade 模式:对外暴露两个语义化方法,隐藏内部 ChainingStrategy 枚举的 4 种状态(ALWAYS, NEVER, HEAD, HEAD_WITH_SOURCES)。你不知道枚举的存在,只需要知道"禁止链"和"开新链"两个概念。这降低了 API 表面的复杂度。
第五关:KeyedStream——只多了两个字段,却解锁了整个有状态世界
打开 KeyedStream.java,跳到第 94-118 行:
publicclassKeyedStream<T, KEY> extendsDataStream<T> {privatefinal KeySelector<T, KEY> keySelector; // ★ key 提取器privatefinal TypeInformation<KEY> keyType; // key 的类型信息publicKeyedStream(DataStream<T> dataStream, KeySelector<T, KEY> keySelector){this(dataStream, keySelector, TypeExtractor.getKeySelectorTypes(keySelector, dataStream.getType())); }}关键细节①:KeyedStream 仍然是一个 DataStream!extends DataStream<T> 意味着 transformation 字段仍然存在,指向父 DataStream 的 Transformation(即 Tokenizer 的 OneInputTransformation(id=2))。keyBy 没有创建新的 DAG 节点。
那么 keySelector 是干嘛的?
告诉下游有状态算子(sum、reduce)如何按 key 分组状态 在创建 PartitionTransformation时,自动附加KeyGroupStreamPartitioner运行时, StreamTask通过它确定当前处理的 key 属于哪个 KeyGroup
为什么 sum 返回 SingleOutputStreamOperator 而不是 KeyedStream? sum 是一个聚合操作,完成后数据已经不需要按 key 分区了。后续的 Sink 不需要状态,不需要 keyBy 的语义。返回 SingleOutputStreamOperator 限制了用户不能再调 keyBy 之后的方法,只能走普通 DataStream 的方法链——这是一种编译期的类型状态机。
DataStream → Transformation 的映射全景
用户代码创建的对象: 内部引用的 Transformation 对象:───────────────────── ──────────────────────────────DataStream<String> text → SourceTransformation(id=1) environment ✓ input: null transformation ✓ outputType: String ↓ text.flatMap(Tokenizer)SingleOutputStreamOperator → OneInputTransformation(id=2)<Tuple2<String,Integer>> input: [SourceTransformation(id=1)]words operatorFactory: SimpleOperatorFactory(StreamFlatMap) environment ✓ outputType: Tuple2<String, Integer> transformation ✓ ↓ words.keyBy(value -> value.f0)KeyedStream → (共享同一个) OneInputTransformation(id=2)<Tuple2<String,Integer>, keySelector 存在 KeyedStream 对象上 String> PartitionTransformation 还没有创建!keyed (只在 sum() 被调用时才创建) keySelector ✓ ↓ keyed.sum(1)SingleOutputStreamOperator → ReduceTransformation(id=4)<Tuple2<String,Integer>> (内部隐含 PartitionTransformation)counts input: [PartitionTransformation → id=2] keySelector: value -> value.f0 operatorFactory: KeyedProcessOperator(ReduceDriver)边界条件速查
DataStream(env, null) | |
DataStream(null, transformation) | |
disableChaining() | |
OutputTag 获取两次但类型不同 |
总结
┌──────────────────────────────────────────────────────────────────┐│ DataStream 类型体系核心要点: ││ ││ 1. DataStream = 壳(environment + transformation),Transformation = DAG 节点 ││ 2. SingleOutputStreamOperator 增加算子配置能力(name, parallelism, chaining) ││ 3. KeyedStream 增加 keySelector,解锁有状态 API(sum, reduce, window) ││ 4. 不同返回类型 = 编译期的方法集控制(类型安全的"状态机") ││ 5. KeyedStream 不创建新 Transformation——分区效果延迟到下游算子创建 ││ 6. Fluent API 全部返回 this 实现链式调用 │└──────────────────────────────────────────────────────────────────┘下一课预告
具体追踪 flatMap() 的完整调用链——FlatMapFunction → StreamFlatMap → SimpleOperatorFactory → OneInputTransformation → 加到 env.transformations List。5 层包装,每一层都有明确的设计理由。
夜雨聆风