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
operator state:在算子层面上持有和管理的状态,不按 key 分割,通常用于在某个任务并行度相同、但要求在不同的并行实例中共享的场景,常见类型包括:ListState
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>