Flink作业中流水线与阻塞数据交换的实现及示例问询
如何创建同时包含流水线与阻塞数据交换的Flink作业
在Flink中,要让作业同时存在**流水线(Pipelined)和阻塞(Blocking)**数据交换,核心是通过控制算子间的数据交换策略,或利用批处理模式下的算子特性实现:
- 流水线交换:数据产生后立即向下游传递,是流处理默认模式,批处理模式下非全局类算子也会默认使用。
- 阻塞交换:算子需处理完所有输入数据后才输出结果,常见于全局排序、全局聚合等需要全量数据计算的场景,也可手动指定交换模式。
代码示例(Java)
以下示例创建了一个混合两种交换模式的批处理作业:
import org.apache.flink.api.common.RuntimeExecutionMode; import org.apache.flink.api.common.functions.MapFunction; import org.apache.flink.api.java.functions.KeySelector; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; public class HybridExchangeJob { public static void main(String[] args) throws Exception { // 1. 创建执行环境,设置为批处理模式 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setRuntimeMode(RuntimeExecutionMode.BATCH); env.setParallelism(2); // 2. 数据源:生成1-100的数字 DataStream<Integer> source = env.fromSequence(1, 100); // 3. 流水线区域1:Source -> Map -> Filter(默认流水线交换) DataStream<Integer> filtered = source .map((MapFunction<Integer, Integer>) num -> num * 2) .filter(num -> num < 150); // 4. 阻塞交换:Filter -> 全局排序(触发阻塞交换,需收集全量数据) DataStream<Integer> sorted = filtered // 全局排序会强制使用阻塞交换,因为要在单分区内完成全量排序 .sortPartition(new KeySelector<Integer, Integer>() { @Override public Integer getKey(Integer value) throws Exception { return value; } }, org.apache.flink.api.common.operators.Order.ASCENDING); // 可选:手动指定阻塞交换(比如在算子间强制设置) // DataStream<Integer> blockedStream = filtered // .setExchangeMode(org.apache.flink.api.common.exchange.ExchangeMode.BLOCKING); // 5. 流水线区域2:Sorted -> Map -> Print(流水线交换) sorted .map((MapFunction<Integer, String>) num -> "Processed: " + num) .print(); env.execute("Hybrid Pipelined-Blocking Job"); } }
关键说明
- 批处理模式配置:通过
env.setRuntimeMode(RuntimeExecutionMode.BATCH)启用批处理模式,这是实现阻塞交换的基础(流模式下默认全流水线,仅可通过setBufferTimeout(-1)实现准阻塞)。 - 自动触发阻塞交换:
sortPartition(全局排序)会自动触发阻塞交换,因为它需要将所有数据收集到单个分区后才能完成排序,必须等上游全量数据处理完毕才会输出结果。 - 手动指定交换模式:如果需要在自定义算子间强制使用阻塞交换,可通过
setExchangeMode(ExchangeMode.BLOCKING)直接指定(示例中注释部分展示了该方式)。 - 区域划分:作业会被划分为3个逻辑区域:
- 流水线区域1:
Source -> Map -> Filter(算子间都是流水线交换,数据实时传递) - 阻塞边界:
Filter到SortPartition(阻塞交换,等待全量输入) - 流水线区域2:
SortPartition -> Map -> Print(恢复流水线交换,数据实时输出)
- 流水线区域1:
验证方式
提交作业后,可通过Flink UI的Job Graph查看:
- 流水线区域内的算子会被合并为同色子图
- 阻塞交换对应的算子间会显示明显边界,且上游算子会先完成所有任务后,下游才开始执行
内容的提问来源于stack exchange,提问作者sunny
相关产品推荐
相关产品推荐

