全文共 4,941 字 预计阅读 15 分钟
bg

flink20问-1

1、Flink 是什么?它与其他流处理框架(如 Spark Streaming)有什么区别?

a、Flink 提供了原生的流计算模型,支持低延迟且高吞吐量的数据处理

b、Flink 采用了“时间与状态”的概念,支持事件时间、处理时间和摄取时间,提供强大的窗口操作和状态管理功能

c、Flink 的高可用性和容错机制更为强大,通过 checkpoint 和 savepoint 可以实现精确一次的语义

d、Flink 的流处理程序可以动态地扩展和缩减资源,以应对变化的工作负载,而 spark streaming 的微批模式则在实时性和灵活性方面不足

spark streaming 主要依赖处理时间,虽然支持事件时间,但是实现和灵活性不如 flink;只能提供至少一次的语义,可以通过一些复杂的设置实现精确一次语义;适用于需要结合批处理和流处理的场景,flink 基于 datastream api 用于流处理,dataset api 用于批处理,spark streaming 基于 rdd api,提供 dstream api 进行流处理

2、在 Flink 中,什么是 DataStream 和 DataSet?两者的主要区别是什么?

datastream 用于处理实时数据流,如实时日志、传感器数据等,使用 datastream 时,数据是动态的,通常需要处理不断到达的新数据;dataset 适用于处理批量数据,如历史数据、批处理任务等,使用 dataset 时,数据是静态的,整个数据集一开始就已经存在,flink 会以微批处理的方式进行处理,datastream 基于 flink 的 checkpoint 来实现容错,dataset 使用较少的容错机制;

3、Flink 的基本架构是什么?包括哪些核心组件?

Flink 的基本架构主要分为以下几个核心组件:JobManager、TaskManager、Job、Execution Graph

JobManager:是 Flink 的控制器,负责处理用户提交的作业,协调作业的执行,并处理任务的调度、资源分配和故障恢复,主要有两个子模块:

a、ResourceManager:负责管理集群中的资源,并与外部资源管理系统(如 yarn、mesos)通信

b、Dispatcher:处理客户端的请求,并提交作业

c、TaskManager:是 Flink 的工作节点,每个节点运行一组 Task,执行实际的数据处理逻辑,负责管理本地资源,包括内存、网络、线程池等

d、Job:是用于编写的 Flink 应用程序,包括数据流和执行逻辑,每个 Job 是一个 dag,节点表示操作(map、filter、window),边表示数据流动

e、Execution Graph:是 Job 的物理计划,在 flink 中负责控制作业执行,将 job 转化为多个并行执行的子任务,并管理这些任务的调度和执行

checkpointing:jobmanager 定期触发检查点操作,将算子状态存储到持久化存储系统中

state management:flink 提供了丰富的状态管理能力,包括 keyed state(键控状态) 和 operator state(算子状态),这是的状态操作变得便捷、高效,同时支持容错。

job execution:flink 支持 batch 和 stream 两种执行模式,在批处理模式下,flink 会读取全部数据后再进行处理,在流处理模式下,数据是连续流入并实时处理的

flinksql:flink 提供了 sql 接口,使得用户可以使用 sql 查询来处理流或批数据

connector 和库:flink 提供了多种连接器,可以与外部系统(如 kafka、hadoop、swa s3等)进行无缝集合

4、在 Flink 中,如何创建一个简单的 DataStream 作业?

a、配置好 flink 运行环境

b、创建 flink 程序的入口,创建 StreamExecutionEnvironment 对象

c、定义输入数据的来源,比如从集合、文件、或者 kafka

d、对输入数据进行处理与转换,可以使用 map、filter、keyBy、timeWindow 等操作

e、指定处理后的数据输出位置,如打印到控制台、写入文件或发送到 kafka

f、执行作业:调用 env.execute("job name") 方法来启动 flink 作业

import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

public class FlinkJob {
    public static void main(String[] args) throws Exception {
        // 1)创建执行环境
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        // 2)定义数据源,从一个简单的集合中读取数据
        DataStream<String> text = env.fromElements(
                "Hello Flink",
                "Flink Stream",
                "Flink DataStream"
        );

        // 3)定义数据转换操作
        DataStream<Integer> wordCounts = text.map(new MapFunction<String, Integer>() {
            @Override
            public Integer map(String value) throws Exception {
                return value.split(" ").length;
            }
        });

        // 4)定义数据输出,输出到控制台
        wordCounts.print();

        // 5)执行作业
        env.execute("Simple DataStream Job");
    }
}

flink 可以通过 setParallelism 方法来指定并行度

5、什么是 Flink 的有状态流处理?与无状态处理有什么区别?

有状态流处理指的是当数据流在处理时,算子的状态会随着输入的数据变化而变化,并将这些状态信息保存在内存或外部存储中,这个状态能在不同时间点被访问并修改,从而维持计算的连续性。例如 flink 可以在流处理过程中维护一个计算器来记录处理过的数据条数,或者保存某些中间计算的结果。

无状态处理指的是每条数据的处理都是相互独立的,不需要依赖于之前的数据或者状态信息;

有状态处理的应用场景:会话窗口(根据时间间隔分割用户的操作记录,识别用户会话)、在线机器学习(算子需要持续更新模型参数,这些参数需要保存在状态中)、数据聚合和累积计算

状态的存储方式:状态可以保存在内存中,速度快但是易受节点故障影响,也可以保存到 hdfs、s3 等,即使节点故障也能恢复状态

flink 使用检查点机制保存状态,并在故障恢复时重建状态,提高系统容错性,通过两阶段提交协议,flink 能保证状态的一致性,避免数据混乱

6、在 Flink 中,如何定义窗口操作?有哪些常见的窗口类型?

窗口操作用于将无界的流式数据划分为有限的片段进行处理,常见的窗口类型包括:时间窗口:按照时间维度划分,如滚动时间窗口和滑动时间窗口;按记录数量划分,如滚动计数窗口和滑动计数窗口;会话窗口:根据一段无数据的静默期划分;

窗口触发器:用于定义何时出发窗口计算和输出,如基于事件数目、基于事件或者自定义条件

窗口处理函数:用于定义窗口计算的逻辑,如 ReduceFunction、AggregateFunction、ProcessWindowFunction

7、Flink 的 Event Time 和 Processing Time 有什么区别?各自的应用场景是什么?

Event Time 指的是事件在其生成源上的时间戳,允许系统根据事件的发生时间来进行处理,从而得到精确的时域计算结果,主要用于需要高精确度时序计算的场景,例如金额交易、实时监控;

Processing Time 指的是事件到达 Flink 系统并被处理的时间,使用 processing time 会根据系统当前的处理时间来进行时间计算,processing time 主要用于那些对时间要求不高或者希望最大化处理性能的场景,例如日志分析、数据清洗

event time 的优点:精确,由于是基于事件发生的时间戳,能够精确地对事件进行排序和窗口计算;缺点:有延迟,需要在事件中携带时间戳信息,并且可能引入延迟,因为需要等待所有数据进入系统才能完成某些计算;

processing time 的优点:简单、高效不用处理时间戳,同样不会引入额外的延迟;缺点:不精确,相同的窗口在不同时次的运行中可能得到不同的结果,缺乏对事件时序的控制;

8、Flink 中的 Watermark 是什么?它在处理延迟数据时的作用是什么?

watermark 是一种特殊的机制,用于处理流式数据中的延迟问题。它标记了数据流中的某个时间点,表示在此时间点之前的数据都已经到达,当 watermark 通过某个操作算子时,flink 会认为晚于 watermark 标记时间的数据已经完成处理;

watermark 可以根据 event time 或者 processing time 生成,通常使用 event time 来生成 watermark,因为它更能反映出来数据生成的真实顺序,可以通过自定义 watermark 生成器来实现,比如实现 AssignerWithPeriodicWatermarks 接口;

watermark 对窗口计算至关重要,flink 会将晚于 watermark 到达的数据视为迟到数据,处理方式是丢弃或策略性输出到侧输出流。

可以设置允许迟到的时间窗,通过调用 window.allowedLateness() 方法实现;

watermark 会沿着数据流向下游传播,当 watermark 到达某个算子,该算子会更新当前 watermark,如果算子有多个输入流,通常会取最小的 watermark 作为当前 watermark

9、Flink 中的 KeyedStream 是什么?它与普通的 DataStream 有什么不同?

KeyedStream 是 Flink 中一种特殊类型的流,通过 key 进行分割,允许基于 key 的状态操作,与普通的 DataStream 不同,KeyedStream 在进行某些操作(如 aggregations 和 windows)时提供了天然的按键分组功能,这可以实现对数据进行精细粒度的处理,如按照用户 ID 汇总购买记录,简单来说 datastream 是未分组的普通数据流,keyedstream 是按 key 分组的 datastream,允许更复杂的状态的窗口操作;

datastream 就是没有内置状态管理机制,更适用 map、flatMap、filter 等无状态操作

// 创建一个基础的 DataStream
DataStream<Tuple2<String, Integer>> dataStream = ...

// 将 DataStream 转换为 KeyedStream
KeyedStream<Tuple2<String, Integer>, String> keyedStream = dataStream
    .keyBy(value -> value.f0); // 按第一个元素(key)进行分组

// 在 KeyedStream 上执行 reduce 操作
SingleOutputStreamOperator<Tuple2<String, Integer>> result = keyedStream
    .reduce((value1, value2) -> new Tuple2<>(value1.f0, value1.f1 + value2.f1));

datastream 适合没有分组需求的 etl;keyedstream 更适合需要按类别处理和分析,如用户行为分析、实时统计、计数操作等

10、在 Flink 中,如何进行状态管理?常见的状态类型有哪些?

flink 主要提供了两种状态类型:Keyed State 和 Operator State,这些状态都可以通过内置的状态后端进行存储和管理,比如内存状态后端、文件系统状态后端和 RocksDB 状态后端

Keyed State 是按键进行管理的状态,通常在 KeyedStream 中使用,每个 key 都可以有自己的状态,这使得状态能被相应的 key 分区隔离和管理,这种状态比较常见的类型有: ValueState:存储一个值;ListState:存储一个列表;MapState<K,V>:存储一个键值对的映射;ReducingState:适用于累积操作,比如求和;AggregatingState<IN, OUT> 适用于自定义的聚合操作。典型应用:点击量统计、各种算子的临时缓存数据

operator state:在算子层面上持有和管理的状态,不按 key 分割,通常用于在某个任务并行度相同、但要求在不同的并行实例中共享的场景,常见类型包括:ListState 适用于保存当前算子所有并行实例共享的一组状态。典型应用:一个 source 算子在使用 FlinkKafkaConsumer 读取数据时,可以将 Kafka 的 offset 信息存储在 Operator State 中,便于断点续读;

Flink 提供了多种状态后端来持久化这些状态:MemoryStateBackend 将状态存储在 TaskManager 的堆内存中,适用于轻量级、小规模的作业;FsStateBackend 将状态存储在文件系统中(如 HDFS),适用于中等规模的作业;RocksDBStateBackend 使用 RocksDB 将状态存储在本地磁盘里,允许存储大量的数据,适用于大规模作业;

flink 通过 checkpointing 和 savepoint 机制来实现状态的持久化和恢复:checkpointing 周期性地拍摄数据流和状态的快照,在失败时自动恢复,面向稳定运行和处理中断;savepoint 是由用户触发的快照,用于升级程序或从指定位置重新启动作业;

flink 的状态管理广泛用于:实时数据分析和统计、实时网络计算、etl 流程

11、Flink 中的 Checkpoint 机制是什么?它如何保证作业的高可用性?

checkpoint 是一种状态管理和恢复机制,用于确保流处理作业的高可用性和一致性,它通过定期将作业中所有算子的状态保存到外部存储,如 hdfs、S3 种,以便在发生故障时能够从最新的 checkpoint 恢复,使得数据处理可以从中断点继续,无需从头开始;

a、在 flink 作业中配置 checkpoint 的参数,如 checkpoint 的间隔时间、存储位置;

b、flink 在运行作业过程中,就定期触发 checkpoint 操作,将当前算子的状态捕获并保存到预设的存储位置

c、一旦作业失败,flink 会自动从最新的 checkpoint 恢复所有算子的状态和数据,使计算捕获并保存到预设的存储位置

d、状态背压:在 checkpoint 过程中数据生产与消费的速度不同步,导致存储或者网络资源的压力,flink 通过异步快照的方法缓解这种状态背压情况,从而使得数据流在进行 checkpoint 时不会中断

e、checkpoint 机制基于两阶段提交协议,首先会发送 checkpoint 触发请求到所有的源,如 kafka 消费者,然后这些源会将当前的状态复制到稳定存储,确认无误后向 flink jobmanager 汇报完成,这样在状态一致的前提下完成 checkpoint

12、在 Flink 中,如何使用算子进行数据转换?有哪些常用的算子?

# map:对每条数据进行转换,返回一条新的数据;
DataStream<String> text = ...;
DataStream<String> uppercased = text.map(s -> s.toUpperCase());

# filter:对数据进行过滤,返回符合条件的数据;
DataStream<String> text = ...;
DataStream<String> filtered = text.filter(s -> s.length() > 3);

# flatMap:将一条数据转换为多条数据;
DataStream<String> sentences = ...;
DataStream<String> words = sentences.flatMap((String sentence, Collector<String> out) -> {
    for (String word : sentence.split(" ")) {
        out.collect(word);
    }
});

# keyBy:对数据按照某个字段进行分组;
DataStream<Tuple2<String, Integer>> userScores = ...;
KeyedStream<Tuple2<String, Integer>, String> keyedByUser = userScores.keyBy(t -> t.f0);

# reduce:对分组后的数据进行聚合;
KeyedStream<Tuple2<String, Integer>, String> keyedByUser = ...;
DataStream<Tuple2<String, Integer>> totalScores = keyedByUser.reduce((value1, value2) -> 
    new Tuple2<>(value1.f0, value1.f1 + value2.f1));

13、Flink 的流处理与批处理是如何统一的?它如何实现“流批一体化”?

flink 的流批一体主要是通过 datastream 和 dataset api 的融合实现的,flink 的底层运行时完全基于流数据,因此批处理作业都是处理有界数据流的特定情况;

flink 一开始有两个 api,datastream 用于处理无界数据,dataset 用于处理有限的数据集,但是随着 flink 的发展,dataset api 的功能被逐步整合到 datastream api 中;

flink 引入了一套 table api 和 sql 来实现 流和批的统一,使开发者能够用声明式的方式进行数据处理;

对于批处理而言,event time 就是数据生成的实际时间,可以通过时间戳字段进行处理

14、在 Flink 中,如何处理数据倾斜问题?有哪些常见的优化策略?

a、数据预处理:在数据进入 flink 之前,对数据进行重新分区,使其分布更加均匀

b、如果某些数据需要在多个任务之间共享,考虑使用广播变量,减少数据重复传递

DataStream<Tuple2<String, Integer>> dataStream = ...;
DataStream<Map<String, String>> broadcastStream = ...;
BroadcastStream<Map<String, String>> bStream = broadcastStream.broadcast();

dataStream.connect(bStream)
    .process(new BroadcastProcessFunction<Tuple2<String, Integer>, Map<String, String>, R>() {
        @Override
        public void processElement(Tuple2<String, Integer> value, ReadOnlyContext ctx, Collector<R> out) {
            // 处理每条数据,结合广播变量
        }

        @Override
        public void processBroadcastElement(Map<String, String> value, Context ctx, Collector<R> out) {
            // 处理广播变量的数据
        }
    });

c、在 keyBy 操作后,随机将数据重新分配到不同的字任务上,以减少某些特定键值导致的数据倾斜

dataStream.keyBy("key")
          .map(new MapFunction<Type, Type>() {
              @Override
              public Type map(Type value) {
                  int randomPartition = new Random().nextInt(numPartitions);
                  return Tuple2.of(randomPartition, value.f1);
              }
          })
          .keyBy(0);

d、将处理流程分为多个阶段,每个阶段使用不同的分区键进行数据重分配

SingleOutputStreamOperator<Tuple2<String, Integer>> phase1 = dataStream.keyBy("key1").sum(1);
SingleOutputStreamOperator<Tuple2<String, Integer>> phase2 = phase1.keyBy("key2").sum(1);

e、自定义分区器,根据数据的特征进行更加细粒度的分区控制

public class CustomPartitioner implements Partitioner<String> {
    @Override
    public int partition(String key, int numPartitions) {
        int partition;
        // 自定义分区逻辑
        return partition;
    }
}

dataStream.partitionCustom(new CustomPartitioner(), "key");

15、在 Flink 中,如何实现窗口的滚动和滑动?

使用 TimeWindows 和 SlidingWindows

滚动窗口:

stream
   .keyBy(<keySelector>)
   .window(TumblingEventTimeWindows.of(Time.minutes(5)))
   .<windowFunction>();

滑动窗口:

stream
   .keyBy(<keySelector>)
   .window(SlidingEventTimeWindows.of(Time.minutes(5), Time.minutes(1)))
   .<windowFunction>();

窗口函数:reduce、fold、apply、aggregate

stream
   .keyBy(<keySelector>)
   .window(<windowType>)
   .reduce(new ReduceFunction<Type>() {
       @Override
       public Type reduce(Type value1, Type value2) {
           // 实现减少逻辑
       }
   });

窗口触发器:触发器定义何时对窗口内容进行计算,可以使用默认触发器,或者自定义触发器来满足特定需求

16、在 Flink 中,如何保证 Exactly Once 语义?它的底层机制是什么?

使用 checkpoints 和 两阶段提交协议。

配置 checkpoints:

env.enableCheckpointing(10000); // 每10秒进行一次checkpoint
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(500); // Checkpoint之间最短间隔500ms
env.getCheckpointConfig().setCheckpointTimeout(60000); // 每次checkpoint超时时间为60s

两阶段提交协议:准备阶段:各个任务不会讲数据提交到外部系统,而是将数据存储在临时位置,并记录一个预提交点;提交阶段:在 checkpoint 成功后,flink 会通知各个任务将先前保存的数据实际提交到外部系统,从而保证数据一致性;

除了确保内部的 exactly once 语义,端到端的 exactly once 还需确保数据源和数据汇也支持,比如可以在 kafka 中配置 producer 的幂等性和 transactions;

flink 有多种状态后端,用于存储作业状态,选择合适的状态后端可以优化性能和稳定性

17、Flink 的 Operator Chain 是如何工作的?如何通过调整链优化作业性能?

Operator Chain 是一个优化数据流系统性能的重要机制,通过在同一个线程中链式地执行一系列操作,可以减少线程间的上下文切换和数据序列化开销,从而提升性能

工作原理:

a、在作业提交阶段,flink 会根据操作符之间的依赖关系来决定哪些操作符可以链在一起

b、被链在一起的操作符会共享同一个输入输出缓冲区,减少了数据复制和传递的开销

c、flink 以拓扑顺序安排操作符的执行,对于链在一起的操作符,这些操作符会在一个线程内依次被执行,不需要进行线程间的切换

调整和优化 Operator Chain 的性能:

a、提高操作符的并行度,可以提高吞吐量

b、通过 disableChaining() 方法显示地控制哪些操作符不被链在一起

c、使用 map 和 filter 操作符进行优化,这些算子通常具有链式优化潜力,可以优先考虑进行链式执行

d、确保数据的分区和重分区策略与链式执行匹配,减少不必要的数据传输

某些场景下,取消链式执行反而会得到更好的性能表现:

a、算子需要更多的资源,cpu / 内存,如某些复杂的函数操作可能导致单个线程资源耗尽、此时拆链可以进行负载均衡

b、在某些情况下,不同的算子链式执行会引入调试和监控上的困难,通过拆链可以分别监控每个算子的性能和资源消耗

c、在链式执行中,如果某个算子出现反压,整个链都可能会收到影响,拆链可以避免这个问题

18、Flink 中的 Side Output 是什么?如何使用 Side Output 实现分流?

side output 是 flink 的一种辅助输出机制,允许在进行流处理时将一部分数据发送到主输出之外的其他流中,这在处理异常数据、不同类型的数据或实现多流操作比较有用;

使用 side output 实现分流的步骤主要包括:

a、定义 side output 标签

b、在 ProcessFunction 或 KeyedProcessFunction 中通过 context.output 方法将数据输出到 side output

c、使用 getSideOutput 方法从 Side Output 中获取数据流

简单例子:

# 定义 OutputTag
OutputTag<String> outputTag = new OutputTag<String>("side-output"){};

# 使用 context.output
public static class MyProcessFunction extends ProcessFunction<String, String> {
    private final OutputTag<String> outputTag;

    public MyProcessFunction(OutputTag<String> outputTag) {
        this.outputTag = outputTag;
    }

    @Override
    public void processElement(String value, Context ctx, Collector<String> out) {
        if (value.contains("error")) {
            ctx.output(outputTag, value);
        } else {
            out.collect(value);
        }
    }
}

# 获取 Side Output 流
DataStream<String> mainStream = /* ... */;
OutputTag<String> outputTag = new OutputTag<String>("side-output"){};

DataStream<String> sideOutputStream = mainStream
    .process(new MyProcessFunction(outputTag))
    .getSideOutput(outputTag);

19、Flink 如何处理有界流和无界流?它们的处理方式有何不同?

有界流:批处理模式:flink 会将数据整体读取到内存,然后进行处理,这种模式可以利用全局视角进行全局优化,比如全局排序、聚合等;无界流:流处理模式,通过时间窗口、滚动窗口等方式进行切分和处理,flink 支持基于事件事件处理,可以处理乱序事件,通过 watermark 等机制,确保正确的时间序处理

20、在 Flink 中,如何通过 Checkpoint 和 Savepoint 实现容错和任务恢复?

checkpoint 由 flink 在后台定期生成 checkpoint,并将其存储在持久化存储中,当任务失败时,flink 会回滚到最近的 checkpoint;savepoint 的生成和恢复过程与 checkpoint 类似,但更为灵活,可以在任务更新或变更时用于恢复

checkpoint 的保存是异步的,不会阻塞数据流的处理,savepoint 一般用户可以通过命令行或者 api 手动触发 savepoint

# 创建 savepoint
./bin/flink savepoint <JobID> [target directory]

# 恢复任务
./bin/flink run -s <SavepointPath> <FlinkJobJar>
Back to Blog