Flink 反压优化

Flink反压优化

在Flink中,反压是指当一个任务(task)生成数据的速率超过下游任务消费数据的速率时触发的警告。以一个简单的Source -> Sink任务链为例,如果Source生成数据的速率高于Sink消费数据的速率,就会发生反压警告。这表示数据流向下游,但由于下游处理速度较慢,系统发出警告以通知上游任务,反压信息沿着相反的方向传播,向上游传递,帮助系统调整任务链以维持平衡,防止性能问题。

Flink任务中出现反压会有如下影响

  1. Flink任务处理性能下降:反压会导致任务链中某些任务被迫减缓数据生成速率,以适应下游任务的处理速度,从而影响整体性能。如:Flink消费Kafka数据时,反压会导致Kafka数据消费有很大滞后。
  2. checkpoint时间长或失败:Flink任务出现反压会使数据处理速度变慢,甚至是数据阻塞,这将导致Checkpoint barrier流经整个数据管道的时间变长,进而Checkpoint的总体时间变长,甚至导致Checkpoint失败。
  3. 内存OOM:在checkpoint barrier对齐下的Exactly-once场景中,由于部分并行度的反压,将导致barrier缓慢到达,数据处理快的并行度会将数据缓存起来以等待数据处理缓慢并行度的barrier对齐,这样可能会导致并行度快的数据一直积压到内存中,导致State占用大量内存资源,最终可能出现OOM问题。
  4. 任务卡住:在下游有窗口的计算逻辑中,如果上游出现数据持续反压,将导致watermark一直不往下游流动,这时导致窗口一直不触发,出现任务卡住情况。

Flink反压问题定位

Flink WebUI提供了Flink 任务运行时的JobGraph图,通过查看该图我们很方便的确定Flink任务是否有反压。在WebUI界面中的JobGraph图中每个算子链默认展示了所有SubTasks的反压和繁忙指标的最大值,如下图所示:

在JobGraph图中除了显示原始的数值外,task也用不同颜色进行了标记,闲置的 tasks 为蓝色,完全被反压的 tasks 为黑色,完全繁忙的 tasks 被标记为红色。 中间的所有值都表示为这三种颜色之间的过渡色。

通过Flink WebUI我们大体能判断Flink Job是否存在反压,但由于Flink默认会将多个算子合并成算子链,通过JobGraph我们只能看到存在反压的算子链,无法具体定位到某个存在反压的算子,这时我们就需要拆成多个步骤来定位具体出现反压的算子,进而定位Flink出现反压的具体原因。

可以按照如下步骤进行Flink任务反压问题定位。

  1. Flink任务禁用算子链后运行Flink任务。
  2. 根据JobGrap定位出现反压的位置。
  3. 结合WebUI Task执行情况定位问题点。
  4. 解决问题点。

下面我们通过一个Flink案例来说明如何一步步定位Flink反压问题,找出出现反压具体原因。有如下存在反压问题的Flink代码:

//1.使用本地模式
Configuration conf = new Configuration();
//设置WebUI绑定的本地端口
conf.setString(RestOptions.BIND_PORT,"8081");
//使用配置
StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(conf);

env.setParallelism(8);

DataStreamSource<StationLog> ds1 = env.addSource(new RichParallelSourceFunction<StationLog>() {
    Boolean flag = true;

    /**
     * 主要方法:启动一个Source,大部分情况下都需要在run方法中实现一个循环产生数据
     * 这里计划1s 产生1条基站数据,由于是并行,当前节点有几个core就会有几条数据
     */
    @Override
    public void run(SourceContext<StationLog> ctx) throws Exception {
        Random random = new Random();
        String[] callTypes = {"fail", "success", "busy", "barring"};
        while (flag) {
            String sid = "sid_" + random.nextInt(10);

            if (sid.equals("sid_0")) {
                for (int i = 0; i < 1000; i++) {
                    String callOut = "1811234" + (random.nextInt(9000) + 1000);
                    String callIn = "1915678" + (random.nextInt(9000) + 1000);
                    String callType = callTypes[random.nextInt(4)];
                    Long callTime = System.currentTimeMillis();
                    Long durations = Long.valueOf(random.nextInt(50) + "");

                    ctx.collect(new StationLog(sid, callOut, callIn, callType, callTime, durations));

                }
                Thread.sleep(200);

            } else {
                String callOut = "1811234" + (random.nextInt(9000) + 1000);
                String callIn = "1915678" + (random.nextInt(9000) + 1000);
                String callType = callTypes[random.nextInt(4)];
                Long callTime = System.currentTimeMillis();
                Long durations = Long.valueOf(random.nextInt(50) + "");

                ctx.collect(new StationLog(sid, callOut, callIn, callType, callTime, durations));
                Thread.sleep(200);
            }

        }

    }

    //当取消对应的Flink任务时被调用
    @Override
    public void cancel() {
        flag = false;
    }
});

//过滤通话状态为fail的数据
SingleOutputStreamOperator<StationLog> ds2 = ds1.filter(new FilterFunction<StationLog>() {
    @Override
    public boolean filter(StationLog value) throws Exception {
        return !"fail".equals(value.callType);
    }
});

//对数据进行keyby
KeyedStream<StationLog, String> ds3 = ds2.keyBy(new KeySelector<StationLog, String>() {
    @Override
    public String getKey(StationLog value) throws Exception {
        return value.sid;
    }
});

//对数据进行聚合操作
SingleOutputStreamOperator<StationLog> ds4 = ds3.sum("duration");

SingleOutputStreamOperator<String> result = ds4.map(new MapFunction<StationLog, String>() {
    @Override
    public String map(StationLog value) throws Exception {
        Thread.sleep(1);
        return "基站ID:" + value.sid + ",通话时长:" + value.duration;
    }
});

result .print();

env.execute();

以上Flink代码中我们通过自定义数据源来产生基站日志数据,产生每条数据的间隔为200ms,当基站为“sid_0”时,立即生产1000条基站为“sid_0”的数据,后续针对这些数据按照基站进行分组,统计每个基站的通话总时长,模拟出数据按照基站分组时存在数据倾斜。在统计得到每个基站总通话时长后,对数据进行转换,处理每条数据时暂停1毫秒。由于数据存在倾斜并且每条数据暂停处理1毫秒会导致下游处理数据的速度小于上游产生数据的速度,进而出现反压现象。

运行以上案例时,我们基于本地运行环境运行,代码中开启了本地WebUI,代码运行后,可以在登录https://localhost:8081进入 Flink WebUI查看Flink任务出现反压。当代码运行一段时间后,Flink Job出现反压如下:

image-20260630104845490

通过上图,我们只能定位到某个算子链中存在反压情况,定位不到反压出现的具体原因,需要按照步骤一步步定位出现反压原因。

  1. 代码设置禁用算子链并重新运行,查看任务反压情况

    在Flink代码中禁用算子链,Flink DataStream和FlinkSQL中禁用算子链方式如下:

    // DataStream 设置方式
    env.disableOperatorChaining();
    
    // FlinkSQL中设置方式
    tableEnv.getConfig().set("pipeline.operator-chaining",false);
    

    设置禁用算子链后,重新运行代码,每隔几秒刷新Flink WebUI可以看到如下反压情况:

    image-20260630111744149

    我们可以看到由于在代码map算子中处理每条数据暂停1毫秒导致Map算子处理数据非常繁忙,进而导致数据按照数据流反方向一层层进行反压。

  2. 根据JobGrap定位出现反压的位置

    在 Flink 中当上游算子在 WebUI 显示有反压时,一般是下游算子存在性能问题,可以继续向下游算子进行排查,直到没有出现反压的算法为止,往往该算子处于繁忙状态,该算子极有可能出现性能问题。

    例如:在以上运行任务中,最开始出现反压的操作为“Keyed Aggregation”,定位反压就找该操作下游算子Map,通过WebUI我们可以看到Map没有出现反压且处于繁忙状态,那么该算子可能存在性能问题。如果Map算子也出现了反压,那就继续向下游进行排查,直到找到没有出现反压的算子为止,说明当前算子出现了性能问题。

    除此外,我们还可以点击每个操作对应的Graph图,查看每个操作对应Task的反压情况,来确定没有反压的第一个操作,如下图所示。

    Source操作对应Task反压图:

    image-20260630112807739

    Filter操作对应Task反压图:

    image-20260630112838259

    KeyBy聚合操作对应Task反压图:

    image-20260630113020575

    Map聚合操作对应Task反压图:

    image-20260630113040181

    以上图中会看到每个操作的反压颜色不同,这些颜色表示意义如下:

    • OK: 0% <= 反压比例 <= 10%
    • LOW: 10% < 反压比例 <= 50%
    • HIGH: 50% < 反压比例 <= 100%

    并且 WebUI 中展示的 Backpressured/Idle/Busy 表示每个task被反压、闲置和繁忙的时间百分比。当刷新该 Flink 任务时,会看到算子数据反压一层层反向阻塞数据,数据流动越来越慢,影响Flink执行效率。

    通过以上方式我们也可以确定Map操作没有数据反压,但对应的1号subtask 非常Busy,可能性能出现问题。

  3. 结合WebUI Task执行情况定位问题点。

    可以进一步查看每个操作对应的SubTask确定具体性能问题点。从Map操作反向查看每个操作SubTask执行情况:

    Map操作Subtask执行情况如下,其中1号task处理数据相比起其他task多很多,出现数据倾斜。

    KeyBy聚合操作Subtask执行情况如下,其中1号task处理数据相比起其他task多很多,出现数据倾斜。

    Filter操作Subtask执行情况如下,其中1号task处理数据相比起其他task多很多,出现倾斜问题。

    根据以上每个操作subtask执行情况,可以大致判断出由于keyby对key进行分组,可能导致数据出现了倾斜,并且由于聚合之后数据默认采用FORWARD(并行分区)方式,数据向下游传递过程中都是1号task出现数据倾斜。参考Map算子BackPressure的情况大致可以判断出由于数据倾斜导致1号task一直处于处理数据状态,进而出现数据反压。

    此外,我们还可以通过火焰图来查看Flink业务逻辑是否出现性能问题。火焰图是一种可视化工具,可以有效地显示Flink中subtask操作占用资源时间长短信息。火焰图是通过多次采样堆栈信息来构建的,每个方法调用由一个柱状图表示,其中柱状图的长度表示某个调用方法在样本中出现的次数多或者少,表示该方法执行时间的长短,长度越长表示执行时间越长。高度由下到上表示Flink执行该方法调用逻辑的顺序。火焰图如下所示:

    图中On-CPU表示线程状态为[RUNNABLE,NEW];Off-CPU表示线程状态为[TIMED_WAITING, WAITING, BLOCKED],Mixed表示两者结合。

    通过火焰图,我们还可以查看具体某个操作中各个subtask执行的CPU占用信息,如下图所示:

    为了避免性能浪费,执行Flink任务时默认没有开启火焰图,我们可以在Flink的conf/ flink-conf.yaml中设置rest.flame .enabled: true或者在提交Flink任务时代码中设置参数开启火焰图即可。

    //开启火焰图
    conf.setString("rest.flamegraph.enabled","true");
    //火焰图刷新周期,默认1分钟
    conf.setString("rest.flamegraph.refresh-interval","10 s");
    

    在以上案例代码中设置以上参数开启火焰图,可以查看Map操作中每个task执行具体占用CPU情况,这里启动程序后需要等待采样一段时间,可以看到对应的Map操作CPU占用时间主要是我们代码中的113行。

    查看代码会发现这里的逻辑处理每条数据时暂停了1ms导致性能问题。

  4. 解决任务反压

    解决以上代码性能问题,我们只需要设置map操作并行度大于前面keyby聚合操作并行度,使两者之间数据分区方式由FORWARD并行分区改变为REBLANCE轮序分区即可,这样每个task处理的数据变得均匀并且将原来大数据量分担到多个task中执行。当然数据倾斜问题如果是代码执行瓶颈问题还需要从源头解决数据倾斜问题。

    ... ...
    SingleOutputStreamOperator<String> result = ds4.map(new MapFunction<StationLog, String>() {
        @Override
        public String map(StationLog value) throws Exception {
            Thread.sleep(1);
            return "基站ID:" + value.sid + ",通话时长:" + value.duration;
        }
    }).setParallelism(16);
    ... ...
    

Flink反压原因及优化

Flink任务中出现反压的原因及优化解决方案如下:

1) 资源设置不合理:当数据源产生数据速度很快,Flink处理数据跟不上时任务很容易产生反压。这种原因主要是Flink处理任务本身性能不足。这种情况可以结合Flink WebUI查看每个task执行和内存使用情况,适当增大Flink任务资源及并行度来解决反压问题。

2) 突发性数据量激增:突发性数据量增加也会导致Flink任务出现反压,例如某个事件窗口内的事件数量突然激增,导致任务无法在规定时间内完成处理。这种反压问题往往需要观察数据量激增频率,偶发生数据激增问题不需额外处理,频繁性数据激增问题最好适当增加Flink任务并行度分摊负载、适当增大窗口大小来解决。

3) 数据倾斜问题:当Flink出现数据倾斜时,会导致处理数据量大的任务会比其他任务处理慢,导致数据在整个流水线中堆积,最终出现数据反压并且还会进一步加剧反压情况。这种情况下需要根据业务数据情况查找数据倾斜问题点并优化解决数据倾斜问题。如:找到数据倾斜key并打散倾斜key。

4) 代码执行效率问题:当Flink代码执行效率慢也会导致Flink反压,代码效率问题可以通过火焰图来了查看确定。代码效率问题原因可能是多方面的,可以根据具体情况进行优化,例如:源头数据读取采用多并行方式、读取外部数据库数据使用异步IO方式、提高算子并行度、避免算子内出现大状态、checkpoint使用不对齐机制、多步骤分散业务避免在一个算子内实现复杂逻辑。

5) DownStream操作延迟:如果Flink作业的下游处理数据能力不足(如:Sink端写入性能差),也会导致Flink任务反压问题,这种情况大多数是由于下游系统的性能瓶颈或者下游交互延迟操作引起。例如:根据下游不同组件来解决优化(如:clickhouse创建分布式表、增大kafka分区、MySQL创建分区表、批量写出);提高Sink端处理数据能力;采用并行写出数据方式(如RichSinkFunction)。

--- 本文结束 The End ---