You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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");
    }
}

关键说明

  1. 批处理模式配置:通过env.setRuntimeMode(RuntimeExecutionMode.BATCH)启用批处理模式,这是实现阻塞交换的基础(流模式下默认全流水线,仅可通过setBufferTimeout(-1)实现准阻塞)。
  2. 自动触发阻塞交换:sortPartition(全局排序)会自动触发阻塞交换,因为它需要将所有数据收集到单个分区后才能完成排序,必须等上游全量数据处理完毕才会输出结果。
  3. 手动指定交换模式:如果需要在自定义算子间强制使用阻塞交换,可通过setExchangeMode(ExchangeMode.BLOCKING)直接指定(示例中注释部分展示了该方式)。
  4. 区域划分:作业会被划分为3个逻辑区域:
    • 流水线区域1:Source -> Map -> Filter(算子间都是流水线交换,数据实时传递)
    • 阻塞边界:Filter到SortPartition(阻塞交换,等待全量输入)
    • 流水线区域2:SortPartition -> Map -> Print(恢复流水线交换,数据实时输出)

验证方式

提交作业后,可通过Flink UI的Job Graph查看:

  • 流水线区域内的算子会被合并为同色子图
  • 阻塞交换对应的算子间会显示明显边界,且上游算子会先完成所有任务后,下游才开始执行

内容的提问来源于stack exchange,提问作者sunny

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.08 13:10:25