核心文件:KeyedStream.java
从哪里调用
上一课 keyBy 返回了 KeyedStream,现在你在上面调 keyed.sum(1)。keyed 的 transformation 指向 Tokenizer(id=2),keySelector 存着 value -> value.f0。
这一课解决什么问题
keyed.sum(1)sum(1) 意味着对 Tuple2 的第 1 个字段(0-indexed,即第二个字段 Integer count)做求和。这一个小小的 sum 背后,隐藏着 Flink 有状态计算的全部秘密。
为什么 sum 是有状态算子而 flatMap 不是?状态存在哪里?故障恢复时状态怎么找回?本课给出答案。
第一关:sum = aggregate = reduce——模板方法模式
KeyedStream.java 第 751-753 行:
public SingleOutputStreamOperator<T> sum(int positionToSum){return aggregate(new SumAggregator<>(positionToSum, getType(), getExecutionConfig()));}SumAggregator 实现了 AggregationFunction<T>,知道如何从 Tuple2<String, Integer> 中取出 position=1 的字段做加法。
跳到第 1004-1006 行:
protected SingleOutputStreamOperator<T> aggregate(AggregationFunction<T> aggregate){return reduce(aggregate).name("Keyed Aggregation");}sum = aggregate = reduce。sum、min、max、minBy、maxBy 都是通过不同的 AggregationFunction 实现 + reduce() 来实现的。
这是模板方法模式:reduce() 定义了"如何做有状态聚合"的模板——创建 Transformation、管理状态、处理输入。具体的聚合逻辑(加法还是取最小值)由传入的 AggregationFunction 决定。增加一种新的聚合类型只需要写一个新的 Function,不需要动 reduce 的代码。
第二关:reduce()——创建 ReduceTransformation
跳到 KeyedStream.java 第 723-740 行:
public SingleOutputStreamOperator<T> reduce(ReduceFunction<T> reducer){ ReduceTransformation<T, KEY> reduce =new ReduceTransformation<>("Keyed Reduce", environment.getParallelism(), transformation, // ← id=2 (Tokenizer 的 OneInputTransformation) clean(reducer), // ← SumAggregator keySelector, // ← value -> value.f0 getKeyType(), // ← Stringfalse);if (isEnableAsyncState) { reduce.enableAsyncState(); } getExecutionEnvironment().addOperator(reduce);returnnew SingleOutputStreamOperator<>(getExecutionEnvironment(), reduce);}注意和 flatMap 的 doTransform() 的对称性——创建 Transformation → addOperator → 返回 DataStream 壳。唯一不同的是 Transformation 类型:ReduceTransformation 代替了 OneInputTransformation。
关键细节①:keySelector 被传入了 ReduceTransformation。上一课讲过,ReduceTransformation 内部检测到 keySelector != null,会自动创建 PartitionTransformation + KeyGroupStreamPartitioner。这就是 keyBy 的"分区效果"终于落地的时刻。
实践环节
在 reduce() 的 getExecutionEnvironment().addOperator(reduce) 处打断点,看此时 transformations List 里已经有哪些节点。应该能看到 Source(id=1)、Tokenizer(id=2)、以及即将加入的 Counter(id=4)。
第三关(关键!):状态是如何引入的——ReduceDriver 的秘密
ReduceTransformation 在运行时创建 ReduceDriver。这是状态诞生的地方:
// ReduceDriver 内部(简化):publicclassReduceDriver<IN, KEY> implementsStatefulStreamTask<IN, IN> {privatetransient ValueState<IN> state; // ★ 这就是状态!就一个字段!publicvoidprocessElement(StreamRecord<IN> element)throws Exception { IN value = element.getValue(); // 新数据,如 ("hello", 1) IN currentValue = state.value(); // 读状态:上一次累加的结果if (currentValue == null) {// 第一次看到这个 key state.update(value); // 状态 = ("hello", 1) output.collect(value); } else {// 不是第一次——做真正的累加 IN newValue = reducer.reduce(currentValue, value); // SumAggregator.reduce(old, new) state.update(newValue); // ★ 更新状态 output.collect(newValue); // 输出累加后的值 ("hello", 6) } }}关键细节②:每个 key 独立维护一份状态。 当处理 "hello" 时,state.value() 返回的是 "hello" 的上一次累加值。当处理 "world" 时,同一个算子实例上的同一个 ValueState 字段会返回 "world" 的上一次累加值。这是 Flink 的状态后端在幕后做的魔法——ValueState<IN> 不是普通的 Java 变量,而是由 StateBackend 管理的、按 key 分片的分布式状态。
状态的生命周期:
算子启动时:从最近 Checkpoint 恢复 ValueState<IN>(首次启动为 null)每处理一条数据:读状态 → 计算 → 写状态 → 输出 Checkpoint 触发时:当前 ValueState序列化后写入 DFS故障恢复时:从 Checkpoint 读回 ValueState,从断点继续
实践环节
在 processElement 方法中打断点,运行 WordCount。第一次看到 key="hello" 时 state.value() 返回 null,第二次再看到 key="hello" 时返回上一次累加的值。这就是"有状态"最直观的演示。
第四关:SumAggregator 的 reduce 逻辑
// SumAggregator.reduce()(简化):public T reduce(T value1, T value2){// value1: 之前累加的结果(来自状态),如 ("hello", 5)// value2: 新到的数据,如 ("hello", 1)if (value1 instanceof Tuple) { Tuple tuple1 = (Tuple) value1; Tuple tuple2 = (Tuple) value2;// 对 position=1 的字段做加法 tuple1.setField( ((Number) tuple1.getField(position)).doubleValue() + ((Number) tuple2.getField(position)).doubleValue(), position);return tuple1; // 返回 ("hello", 6) }// ...处理其他类型}逻辑不复杂——取出指定位置字段,做数值加法。但要注意一个细节:这里用的是 doubleValue() 做累加,即使原始类型是 Integer。这会导致 sum(1) 的结果字段变成 Double 类型。如果你期望保持 Integer,需要用 reduce(new SumFunction<Integer>()) 自己实现。
有状态 vs 无状态算子对比
OneInputTransformation | ReduceTransformation | |
| 必须 | ||
调用链追踪
keyed.sum(1) │ ├─→ new SumAggregator<>(1, type, config) │ └→ position=1 意味着对 Tuple2 的第2个字段(Integer)做加法 │ ├─→ aggregate(sumAggregator) │ └→ reduce(sumAggregator) // sum = reduce │ └─→ reduce(sumAggregator) │ ├─→ new ReduceTransformation<>( │ "Keyed Reduce", │ env.getParallelism(), │ transformation, // id=2 (Tokenizer) │ sumAggregator, // ReduceFunction │ keySelector, // value -> value.f0 │ keyType, // String │ false) │ │ │ └→ 内部:因为 keySelector != null │ 自动创建 PartitionTransformation( │ Tokenizer(id=2), │ KeyGroupStreamPartitioner(keySelector)) │ ├─→ env.addOperator(reduceTransformation) │ └→ transformations.add(reduce) // 加入 DAG │ └─→ new SingleOutputStreamOperator<>(env, reduce) // 返回给用户总结
┌─────────────────────────────────────────────────────────────────┐│ sum() 核心要点: ││ ││ 1. sum = aggregate = reduce(模板方法模式) ││ 2. SumAggregator 决定"怎么聚合",reduce() 决定"怎么管理状态" ││ 3. ReduceTransformation 在运行时创建 ReduceDriver ││ 4. ReduceDriver 内部维护 ValueState<IN>(本次累加结果) ││ 5. 每个 key 独立维护一份状态("hello"的累加值 ≠ "world"的) ││ 6. 状态随 Checkpoint 持久化,故障恢复时从 Checkpoint 恢复 ││ 7. 必须 keyBy 后才能 sum——状态按 key 组织 │└─────────────────────────────────────────────────────────────────┘下一课预告
阶段二总结——画出 WordCount 的完整 Transformation DAG,理解 5 个节点如何通过 input 引用串联成单向链表。以及这套链式 DAG 如何衔接到阶段三的 StreamGraph。
夜雨聆风