50道大数据精选面试题
Hive 面试题
面试题1:Hive 中四个 by 的区别?
-
Sort By:分区内有序;
- 不是全局排序,其在数据进入 reducer 前完成排序,也就是说它会在数据进入 reduce 之前为每个 reducer 都产生一个排序后的文件。因此,如果用 sort by 进行排序,并且设置
mapreduce.job.reduces > 1,则 sort by 只保证每个 reducer 的输出有序,不保证全局有序。
- 不是全局排序,其在数据进入 reducer 前完成排序,也就是说它会在数据进入 reduce 之前为每个 reducer 都产生一个排序后的文件。因此,如果用 sort by 进行排序,并且设置
-
Order By:全局排序,只有一个 Reducer;
- order by 会对输入做全局排序,因此只有一个 reducer(多个 reducer 无法保证全局有序),然而只有一个 reducer,会导致当输入规模较大时,消耗较长的计算时间。
-
Distribute By:类似 MR 中 Partition,进行分区,结合 sort by 使用。
- distribute by 是控制在 map 端如何拆分数据给 reduce 端的。类似于 MapReduce 中分区 partitioner 对数据进行分区。
- Hive 会根据 distribute by 后面列,将数据分发给对应的 reducer,默认是采用 hash 算法 + 取余数的方式。
-
Cluster By:等同于 Distribute by + Sort by,只能按默认升序排序。
- 当 Distribute by 和 Sorts by 字段相同时,可以使用 Cluster by 方式。Cluster by 除了具有 Distribute by 的功能外还兼具 Sort by 的功能。但是排序只能是升序排序,不能指定排序规则为 ASC 或者 DESC。
面试题2:Hive 中静态分区和动态分区的区别?
静态分区与动态分区的主要区别在于静态分区是手动指定,而动态分区是通过数据来进行判断。
详细来说,静态分区的列是在编译时期,通过用户传递来决定的;动态分区只有在 SQL 执行时才能决定。静态分区不管有没有数据都将会创建该分区,动态分区是有结果集将创建,否则不创建。
-
静态分区 SP(static partition)
- 静态分区是在编译期间指定的指定分区名。
- 支持 load 和 insert 两种插入方式。
- load 方式:
- 会将分区字段的值全部修改为指定的内容。
- 一般是确定该分区内容是一致的时候才会使用。
- insert 方式:
- 必须先将数据放在一个没有设置分区的普通表中。
- 该方式可以在一个分区内存储一个范围的内容。
- 从普通表中选出的字段不能包含分区字段。
- load 方式:
- 适用于分区数少,分区名可以明确的数据。
-
动态分区 DP(dynamic partition)
- 根据分区字段的实际值,动态进行分区。
- 是在 SQL 执行的时候进行分区。
- 需要先将动态分区设置打开(
set hive.exec.dynamic.partition.mode=nonstrict)。 - 只能用 insert 方式。
面试题3:Hive 内部表与外部表的区别?使用场景分别是?
-
区别:
- 内部表数据由 Hive 自身管理,外部表数据由 HDFS 管理。
- 内部表数据存储的位置是
hive.metastore.warehouse.dir。Hive 自身管理,外部表数据由 HDFS 管理。 - 删除内部表会直接删除元数据(metadata)及存储数据;删除外部表仅仅会删除元数据,HDFS 上的文件并不会被删除。
- 对内部表的修改会将修改直接同步给元数据,而对外部表的表结构和分区进行修改,则需要修复。
-
使用场景:
- 外部表:相对来说更加安全些,数据组织也更加灵活,方便共享源数据。如果数据的处理由 Hive 和其他工具一起处理,则创建外部表。
- 内部表:如果所有的数据都由 Hive 处理,则创建内部表。
Kafka 面试题
面试题4:Kafka 如何保证数据不丢失?
Broker 端:
- Topic 副本因子个数:
replication.factor >= 3。 - 同步副本列表(ISR):
min.insync.replicas = 2。 - 禁用 unclean 选举:
unclean.leader.election.enable=false。
Producer 端:
- 同步方式:
producer.type=sync。request.required.acks=1。- 副本数量
>=2。 - 增加重试次数。
- 异步方式:
producer.type=async。request.required.acks=1。queue.buffering.max.ms=5000。queue.buffering.max.messages=10000。queue.enqueue.timeout.ms = -1。batch.num.messages=200。queue.buffering.max.ms=5000。- 通过 buffer 来进行控制数据的发送,有两个值来进行控制,缓冲时间阈值与缓冲消息的数量阈值,如果 buffer 满了数据还没有发送出去,如果设置的是立即清理模式,风险很大,一定要设置为阻塞模式。
Consumer 端:
- 关闭自动 offset,手动提交 offset。
- 设置
enable.auto.commit = false,默认值 true,自动提交。 - 使用 Kafka 的 Consumer 的类,用方法
consumer.commitSync()提交。 - 或者使用 spring-kafka 的 Acknowledgment 类,用方法
ack.acknowledge()提交(推荐使用)。
- 设置
- 另一个方法同样需要手动 commit offset,另外在 consumer 端再将所有 fetch 到的数据缓存到 queue 里,当把 queue 里所有的数据处理完之后,再批量提交 offset,这样就能保证只有处理完的数据才被 commit。
面试题5:Kafka 如何保证数据 exactly-once?
-
Producer exactly-once:
enable.idempotence=true。- 分区副本数
>= 2。 isr >=2。ProducerID + SequenceNumber + Ack = -1(幂等性)。
-
Consumer exactly-once:
- 手动维护并提交偏移量。
- 设置
enable.auto.commit=false,关闭自动提交偏移量。 - 借助外部数据库,如 Redis 的 pipeline,MySQL 的事务机制管理存储偏移量。
- 在同一事物中,在消息被处理完之后再提交偏移量。并更新偏移量。否则消息需回滚,并获取到上一次偏移量的位置重新进行处理。
面试题6:Kafka 数据积压怎么解决?
- 增加 broker 节点,增加分区数量,提高并行度。
- 修改单个消费为批量消费。
- 增加单线程消费为线程池异步消费。
- 缩短批次时间间隔。
- 老版本 SparkStreaming 控制消费的速率(
spark.streaming.kafka.maxRatePerPartition),可以控制最大的消费速率,在参数中设置;新版本设置背压机制实现消费处理的动态平衡。 - 对代码进行优化,尽可能的一次性计算多个结果,减少 shuffle 过程。
- 处理的结果如果过多,可以将数据保存到 MySQL 集群、MongoDB 集群【支持事物】或 ES【不支持事物】,增大吞吐量。
- 消费线程将拉取的消息放到一个滑动窗口中,通过滑动窗口控制拉取的速度。
- 对于倾斜的 key 加以处理,加随机数等方式打散。
面试题7:Kafka rebalance 发生时机和分区分配策略?
Rebalance 触发时机?
- 当出现以下几种情况时,Kafka 会进行一次重新分区分配操作,即 Kafka 消费者端的 Rebalance 操作:
- 同一个 consumer 消费者组
group.id中,新增了消费者进来,会执行 Rebalance 操作。 - 消费者离开当期所属的 consumer group 组。比如主动停机或者宕机。
- 分区数量发生变化时(即 topic 的分区数量发生变化时)。
- 消费者主动取消订阅。
- 同一个 consumer 消费者组
Rebalance 三种策略?
- Kafka 新版本提供了三种 rebalance 分区分配策略:(
partition.assignment.strategy):- range。
- round-robin。
- sticky。
面试题8:Kafka 的分区数如何确定?
在 partition 级别上达到均衡负载是实现吞吐量的关键,合适的 partition 数量可以达到高度并行读写和负载均衡的目的,需要根据每个分区的生产者和消费者的目标吞吐量进行估计。
可以遵循一定的步骤来确定分区数:根据某个 topic 日常“接收”的数据量等经验确定分区的初始值,然后测试这个 topic 的 producer 吞吐量和 consumer 吞吐量。假设它们的值分别是 Tp 和 Tc,单位可以是 MB/s。然后假设总的目标吞吐量是 Tt,那么:
numPartitions = Tt / max(Tp, Tc)
说明:Tp 表示 producer 的吞吐量。测试 producer 通常是很容易的,因为它的逻辑非常简单,就是直接发送消息到 Kafka 就好了。Tc 表示 consumer 的吞吐量。测试 Tc 通常与应用消费消息后进行什么处理的关系更大,相对复杂一些。
分区数过多的危害?
- 客户端/服务器端需要使用的内存就越多。
- 文件句柄的开销。
- 越多的分区可能增加端对端的延迟。
- 降低高可用性。
面试题9:数据发往 Kafka 的分区规则?
key 和 value 的类型,一般都用字符串即可。数据到底写入到哪一个分区中:
- 如果指定了分区,就写入到指定的分区中。
- 如果没有指定分区,指定了 key,按照 key 的 hashcode,取模,写入对应的分区。
- 如果没有指定分区和 key,轮询机制。
面试题10:Kafka producer buffer pool 的作用?
Kafka 通过使用内存缓冲池的设计,让整个发送过程中的存储空间循环利用,有效减少 JVM GC 造成的影响,从而提高发送性能,提升吞吐量。
面试题11:Kafka 时间轮的作用?
Kafka 通过时间轮来处理延迟任务,只将时间轮的槽保存到延迟队列,大大的减少了延迟队列的元素数量,这样对于元素的增加删除性能有很大提高;Kafka 通过阻塞的方式 poll 延迟队列的,减少了大量的空转;为了保证线程安全,灵活运用读写锁、原子对象、synchronized 控制时间轮的操作。
面试题12:Kafka 为什么这么快?

Spark 面试问题
面试题13:Spark 为什么比 MapReduce 快?
这是一道常见的面试题,回答时可以从 IO、shuffle 与排序、资源、部署模式、内存管理策略等各个方面来回答。
- MR 基于磁盘的分布式计算引擎,频繁的磁盘 IO。Spark 基于内存进行计算,DAG 计算模型,大大减少了磁盘 IO。
- Spark 多线程运行,MR 多进程运行。
- Spark 粗粒度资源申请,MR 细粒度资源申请。
- Spark 支持多种部署模式,MR 只支持 yarn 上部署。
- Shuffle 与排序;MR 有 reducer 必排序,一般会经过 3 次排序,Spark Shuffle 数据的排序操作不是必须的。
- Spark 有多种 shuffle 类型,Spark 不一定会发生 shuffle,MR 一定会发生 shuffle。
- Spark 具有灵活的内存管理策略。
面试题14:Spark Repartition 和 Coalesce 的关系与区别,能简单说说吗?
-
关系:
- 两者都是用来改变 RDD 的 partition 数量的,repartition 底层调用的就是 coalesce 方法:
coalesce(numPartitions, shuffle = true)。
- 两者都是用来改变 RDD 的 partition 数量的,repartition 底层调用的就是 coalesce 方法:
-
区别:
- repartition 一定会发生 shuffle,coalesce 根据传入的参数来判断是否发生 shuffle。
- 一般情况下增大 RDD 的 partition 数量使用 repartition,减少 partition 数量时使用 coalesce。
面试题15:简述下 Spark 中的缓存(cache 和 persist)与 checkpoint 机制,并指出两者的区别和联系?
关于 Spark 缓存和检查点的区别,大致可以从这 4 个角度去回答:
- 位置:
- Persist 和 Cache 将数据保存在内存,Checkpoint 将数据保存在 HDFS。
- 生命周期:
- Persist 和 Cache 程序结束后会被清除或手动调用 unpersist 方法,Checkpoint 永久存储不会被删除。
- RDD 依赖关系:
- Persist 和 Cache,不会丢掉 RDD 间的依赖链/依赖关系,CheckPoint 会斩断依赖链。
- 执行与使用:
- persist 中 RDD 的逻辑只会执行一次,而 checkpoint 会执行两次。
- 生产环境中一般都是 cache 和 checkpoint 连用,这样 RDD 逻辑只会执行一次,并且会缓存到 checkpoint 中。

面试题16:Spark on Yarn client 模式与 cluster 的区别?
- driver 所在位置不同:
- client 模式下 driver 线程只在 spark-submit 命令提交的机器上。
- cluster 模式下,driver 线程只在 applicationMaster 所在的节点。
- 启动的任务进程名字不一样:
- client 模式下:ExecutorLauncher 只负责向 yarn 申请容器来启动 executor。
- cluster 模式下,applicationMaster 既要负责申请运行 executor 的资源,又要调 Driver 线程来做 task 调度。
面试题17:Spark RDD、DataFrame、Dataset 的区别与联系?
三者的共性:
- RDD、DataFrame、DataSet 全都是 Spark 平台下的分布式弹性数据集,为处理超大型数据提供便利。
- 三者都有惰性机制,在进行创建、转换,如 map 方法时,不会立即执行,只有在遇到 Action 如 foreach 时,三者才会开始遍历运算。
- 三者有许多共同的函数,如 filter,排序等。
- 在对 DataFrame 和 Dataset 进行操作许多操作都需要这个包:
import spark.implicits._(在创建好 SparkSession 对象后尽量直接导入)。 - 三者都会根据 Spark 的内存情况自动缓存运算,这样即使数据量很大,也不用担心会内存溢出。
- 三者都有 partition 的概念。
- DataFrame 和 Dataset 均可使用模式匹配获取各个字段的值和类型。
三者的区别:
- RDD:
- RDD 一般和 Spark MLlib 同时使用。
- RDD 不支持 SparkSQL 操作。
- DataFrame:
- 与 RDD 和 DataSet 不同,DataFrame 每一行的类型固定为 Row,每一列的值没法直接访问,只有通过解析才能获取各个字段的值。
- DataFrame 与 DataSet 一般不与 Spark MLlib 同时使用。
- DataFrame 与 DataSet 均支持 SparkSQL 的操作,比如 select,groupby 之类,还能注册临时表/视窗,进行 SQL 语句操作。
- DataFrame 与 DataSet 支持一些特别方便的保存方式,比如保存成 csv,可以带上表头,这样每一列的字段名一目了然。
- DataSet:
- DataSet 和 DataFrame 拥有完全相同的成员函数,区别只是每一行的数据类型不同。DataFrame 其实就是 DataSet 的一个特例。
type DataFrame = DataSet[Row]。 - DataFrame 也可以叫 DataSet[Row],每一行类型是 Row,不解析,每一行究竟有哪些字段,各个字段又是什么类型都无从得知,只能用上面的 getAs 方法或者共性中的第七条提到的模式匹配拿出特定字段,而 DataSet 中,每一行是什么类型是不一定的,在自定义 case class 之后可以很自由的获取每一行的信息。
- DataSet 和 DataFrame 拥有完全相同的成员函数,区别只是每一行的数据类型不同。DataFrame 其实就是 DataSet 的一个特例。
三者的转换:
- RDD 转 DataFrame:
- 方案一:直接将字段名称传入 toDF 中:
.toDF(col1, col2...)。 - 方案二:通过反射的方式:
- Java:
.createDataFrame(RDD1, JavaBean.class)。 - Scala:通过 case class 如:
scala
spark.sparkContext.textFile(path).map(line => line.split(",")).map(x => { Person(x(0), x(1).trim.toLong) }).toDF()
- Java:
- 方案三:构造 Schema 的方式:
.createDataFrame(rdd, scheme)。
- 方案一:直接将字段名称传入 toDF 中:
- DataFrame 转 RDD:
df.rdd或者df.javaRDD()。
- RDD 转 DataSet:
- 方案一:使用 toDS() 算子,需要导入隐式转换(
import spark.implicits._)。 - 方案二:使用
spark.createDataset(rdd)。
- 方案一:使用 toDS() 算子,需要导入隐式转换(
- DataSet 转 RDD:
- 直接使用
.rdd。
- 直接使用
- DataFrame 转 DataSet:
- 封装样例类,调用
df.as[xxx]。 case class xxx()。df.as[xxx]。
- 封装样例类,调用
- DataSet 转 DataFrame:
ds.toDF()。

面试题18:updateStateByKey 与 mapWithState 使用区别?
- updateStateByKey:统计全局的 key 的状态,就算没有数据输入,它也会在每一个批次的时候返回之前的 key 的状态。
- 缺点:若数据量太大的话,需要 checkpoint 的数据会占用较大的存储,效率低下。
- mapWithState:也是用于全局统计 key 的状态,但是它如果没有数据输入,便不会返回之前的 key 的状态,有一点增量的感觉。效率更高,生产中建议使用。
- 优点:我们可以只是关心那些已经发生的变化的 key,对于没有数据输入,则不会返回那些没有变化的 key 的数据。这样的话,即使数据量很大,checkpoint 也不会像 updateStateByKey 那样,占用太多的存储。
面试题19:Spark SQL 三种 join 方式?
- Broadcast Hash Join:适合一张很小的表和一张大表进行 Join。
- Shuffle Hash Join:适合一张小表(比上一个大一点)和一张大表进行 Join。
- Sort Merge Join:适合两张大表进行 Join。
Shuffle Hash Join 策略必须满足以下条件:
- 仅支持等值 Join,不要求参与 Join 的 Keys 可排序(这点是和 sort-merge join 相对应)。
spark.sql.join.preferSortMergeJoin参数必须设置为 false,参数是从 Spark 2.0.0 版本引入的,默认值为 true,也就是默认情况下选择 Sort Merge Join。- 小表的大小(
plan.stats.sizeInBytes)必须小于spark.sql.autoBroadcastJoinThreshold * spark.sql.shuffle.partitions(默认值 200)其实就是让每一个小表的分区都类似于广播变量的小表。 - 而且小表大小(
stats.sizeInBytes)的三倍必须小于等于大表的大小(stats.sizeInBytes),也就是a.stats.sizeInBytes * 3 <= b.stats.sizeInBytes。
Broadcast Hash Join 策略必须满足以下条件:
- 小表的数据必须很小,可以通过
spark.sql.autoBroadcastJoinThreshold参数来配置,默认是 10MB。 - 如果内存比较大,可以将阈值适当加大。
- 将
spark.sql.autoBroadcastJoinThreshold参数设置为 -1,可以关闭这种连接方式。 - 只能用于等值 Join,不要求参与 Join 的 keys 可排序。
要启用 Shuffle Sort Merge Join 必须满足的条件是仅支持等值 Join,并且要求参与 Join 的 Keys 可排序。
面试题20:RDD 有什么缺陷?
- 不支持细粒度的写和更新操作(如网络爬虫),Spark 写数据是粗粒度的。所谓粗粒度,就是批量写入数据,为了提高效率。但是读数据是细粒度的,也就是说可以一条条的读。
- 不支持增量迭代计算,Flink 支持。
面试题21:groupByKey 和 reduceByKey 区别?
reduceByKey 和 groupByKey 都存在 shuffle 的操作,但是 reduceByKey 可以在 shuffle 前对分区内相同 key 的数据进行预聚合(combine)功能,这样会减少落盘的数据量,而 groupByKey 只是进行分组,不存在数据量减少的问题,reduceByKey 性能比较高。从功能的角度:reduceByKey 其实包含分组和聚合的功能。GroupByKey 只能分组,不能聚合,所以在分组聚合的场合下,推荐使用 reduceByKey,如果仅仅是分组而不需要聚合。那么还是只能使用 groupByKey。
面试题22:RDD 的弹性表现在哪几点?
- 自动的进行内存和磁盘的存储切换。
- 基于 Lineage 的高效容错。
- Task 如果失败会自动进行特定次数的重试。
- Stage 如果失败会自动进行特定次数的重试,而且只会计算失败的分片。
- Checkpoint 和 persist,数据计算之后持久化缓存。
- 数据调度弹性,DAG task 调度和资源无关。
- 数据分片的高度弹性。
面试题23:RDD 通过 Lineage(记录数据更新)的方式为何很高效?
- Lazy 记录了数据的来源,RDD 是不可变的,且是 lazy 级别的,且 RDD 之间构成了链条,lazy 是弹性的基石。由于 RDD 不可变,所以每次操作就产生新的 RDD,不存在全局修改的问题,控制难度下降,所有有计算链条将复杂计算链条存储下来,计算的时候从后往前回溯 900 步是上一个 stage 的结束,要么就 checkpoint。
- 记录原数据,是每次修改都记录,代价很大如果修改一个集合,代价就很小,官方说 RDD 是粗粒度的操作,是为了效率,为了简化,每次都是操作数据集合,写或者修改操作,都是基于集合的 RDD 的写操作是粗粒度的,RDD 的读操作既可以是粗粒度的也可以是细粒度,读可以读其中的一条条的记录。
- 简化复杂度,是高效率的一方面,写的粗粒度限制了使用场景如网络爬虫,现实世界中,大多数写是粗粒度的场景。
面试题24:Spark 3.0 AQE 新特性?
- 自动分区合并。
- 自动数据倾斜处理。
- Join 策略调整。
- 详细见:Spark 3.0 AQE 新特性。
面试题25:Spark Hash shuffle 与 Sort shuffle 的区别?
Hash shuffle:
- 一种是普通运行机制,另一种是合并的运行机制。
- 产生的磁盘小文件的个数为
maptask * reducetask,每个分区是一个 task,磁盘小文件多,I/O 增多,产生的 GC 会增多。 - 这种 shuffle 产生的磁盘小文件,容易导致 OOM。
- 这种模式不单单产生的磁盘小文件比较多,而且占用内存也比较多。我们应该降低这种磁盘之间的接触。
Hash shuffle 的优化机制:
- 启动 HashShuffle 的合并机制 ConsolidatedShuffle 的配置:
spark.shuffle.consolidateFiles=true。 - 两个 task 共用一个 buffer 缓冲区。
- 如果 Reducer 端的并行任务或者是数据分片过多的话则
Core * Reducer Task依旧过大,也会产生很多小文件。
Sort shuffle:
- Spark 1.6 之前用 hash shuffle,在 Spark 1.6 之后使用 sort shuffle。
- Sort shuffle 的两种机制:
- 估算,去要内存
5.01 * 2 - 5,要不到的时候就去排序,最终溢写的小的磁盘小文件合并成为了一个大的磁盘小文件。 - 当不需要排序的时候,默认使用 Bypass 机制。
- 估算,去要内存
- Bypass 运行机制的触发条件:
- Shuffle reduce task 数量小于
spark.shuffle.sort.bypassMergeThreshold参数的值小于 200,不开启,溢写磁盘不需要排序,小于等于的时候是开启的。 - 不是聚合类的 shuffle 算子(比如 reduceByKey)。
- Shuffle reduce task 数量小于
总结:
- Hash shuffle(合并运行机制)优化机制产生的磁盘小文件的个数:
C * R(core * reducer)。 - Hash shuffle(普通):产生的磁盘小文件:
M * R。 - Sort shuffle 产生的磁盘小文件的个数为:
2 * M。 - Bypass 机制产生的磁盘小文件的个数为:
2 * M。
面试题26:哪些 Spark 算子会有 shuffle 过程?
- 去重:
distinct。 - 排序:
groupByKey,reduceByKey,sortByKey。 - 重分区:
repartition,repartitionAndSortWithinPartitions,coalesce。 - 集合或者表连接操作:
join,cogroup。
Flink 面试题
面试题27:Flink 中什么时候形成算子链,什么时候断链?
形成算子链的条件:
- 上下游的并行度一致(槽一致)。
- 该节点必须要有上游节点跟下游节点。
- 下游 StreamNode 的输入 StreamEdge 只能有一个。
- 上下游节点都在同一个 slot group 中(下面会解释 slot group)。
- 下游节点的 chain 策略为 ALWAYS(可以与上下游链接,map、flatmap、filter 等默认是 ALWAYS)。
- 上游节点的 chain 策略为 ALWAYS 或 HEAD(只能与下游链接,不能与上游链接,Source 默认是 HEAD)。
- 上下游算子之间没有数据 shuffle(数据分区方式是 forward)。
- 用户没有禁用 chain。
禁用算子链的场景:
- 某个算子需要单独设置资源:当某个算子需要单独设置资源时,比如说 Memory、CPU 等,这个算子就不能被放置在算子链里面,需要单独成为一个 Task。(背压时定位问题)。
- 需要等待外部事件触发:某些算子需要等待外部事件触发才能继续处理数据,例如读取外部文件或者接收网络消息等,这时候算子如果被放到算子链里面,则整个链都会阻塞,产生性能问题,因此需要禁用算子链。
- 处理时间窗口非常大的数据集:对于非常大的数据集,特别是在窗口结束时间很大的情况下,算子链可能会消耗太长的时间,导致超时或者 OOM 错误。禁用算子链可以避免此类问题。
- 需要流数据处理与批数据处理共存:如果同时需要进行流式与批处理,禁用算子链可以让流和批处理同时运行,避免出现串行化的问题。
面试题28:Flink keyby 之后出现数据倾斜的原因是什么?如何定位和解决的?
Keyby 之后出现数据倾斜常见原因?
- Key 的选择不合适:如果选择的 Key 不平衡或者有明显的热点数据,就容易出现数据倾斜的问题。应该尝试选择更加平衡的 Key,例如多个属性组合的方式。
- 数据分布不均匀:有些数据在时间、空间上分布不均匀,导致某些 Key 的数据量比其他 Key 大很多。可以通过统计每个 Key 对应的数据量,找到数据分布不均匀的原因。
- 算子链长/复杂度高:当算子链过长或者算子的操作很复杂时,也容易导致某些 Task 的数据处理量过大。可以通过拆分算子链、优化算子操作等方式来解决。
- 并行度设置不当:并行度过高可能导致资源浪费,过低则会导致数据倾斜。应该根据实际情况,合理设置并行度。
定位:
- 步骤 1:定位反压:
- 定位反压有 2 种方式:Flink Web UI 自带的反压监控(直接方式)、Flink Task Metrics(间接方式)。通过监控反压的信息,可以获取到数据处理瓶颈的 Subtask。
- 步骤 2:确定数据倾斜:
- Flink Web UI 自带 Subtask 接收和发送的数据量。当 Subtasks 之间处理的数据量有较大的差距,则该 Subtask 出现数据倾斜。如下图所示,红框内的 Subtask 出现数据热点。
解决方案:
- keyBy 后聚合操作存在数据倾斜(通过 Flink LocalKeyBy 思想来解决):
- 在 keyBy 上游算子数据发送之前,首先在上游算子的本地对数据进行聚合后再发送到下游,使下游接收到的数据量大大减少,从而使得 keyBy 之后的聚合操作不再是任务的瓶颈。类似 MapReduce 中 Combiner 的思想,但是这要求聚合操作必须是多条数据或者一批数据才能聚合,单条数据没有办法通过聚合来减少数据量。从 Flink LocalKeyBy 实现原理来讲,必然会存在一个积攒批次的过程,在上游算子中必须攒够一定的数据量,对这些数据聚合后再发送到下游。
- 注意:Flink 是实时流处理,如果 keyby 之后的聚合操作存在数据倾斜,且没有开窗口的情况下,简单的认为使用两阶段聚合,是不能解决问题的。因为这个时候 Flink 是来一条处理一条,且向下游发送一条结果,对于原来 keyby 的维度(第二阶段聚合)来讲,数据量并没有减少,且结果重复计算(非 FlinkSQL,未使用回撤流)。
- keyBy 后窗口聚合操作存在数据倾斜(两阶段聚合):
- 因为使用了窗口,变成了有界数据的处理,窗口默认是触发时才会输出一条结果发往下游,所以可以使用两阶段聚合的方式:
- 第一阶段聚合:key 拼接随机数前缀或后缀,进行 keyby、开窗、聚合。
- 第二阶段聚合:去掉随机数前缀或后缀,按照原来的 key 及 windowEnd 作 keyby、聚合。
- 因为使用了窗口,变成了有界数据的处理,窗口默认是触发时才会输出一条结果发往下游,所以可以使用两阶段聚合的方式:

面试题29:Flink TaskManager slot JobManager 并行度之间资源怎么分配的?
TaskManager、slot、并行度之间的关系:
- 在 Yarn 集群中 Job 分离模式下,TaskManager 的数量 =
ceil(slot 数量 / 并行度)。slotNumber >= taskManager * 并行度。
TaskManager/slots 与 CPU 的关系:
- 经验上讲 Slot 的数量与 CPU-core 的数量一致为好。但考虑到超线程,可以让
slotNumber = 2 * cpuCore。
slot 与并行度:
- 一般我们设置 task 的并行度不能超过 slot 的数量。
- 一个 Task 的并行度等于分配给它的 Slot 个数(前提槽资源充足)。
面试题30:Flink application、session、pre-job 模式的区别和使用场景?
- application:每个 job 独享一个集群,job 退出则集群退出。main 方法在集群上运行。
- session:多个 job 共享集群资源,job 退出集群也不会退出。main 方法在客户端运行。
- pre-job:每个 job 独享一个集群,job 退出则集群退出。main 方法在客户端运行。
适用场景:
- Session 模式:一般用来部署那些对延迟非常敏感但运行时长较短的作业,需要频繁提交小 job 的场景。
- Per-Job 模式:一般用来部署那些长时间运行的作业。
- Application 模式:综合了两种模式的所有优点,建议生产上适用。

面试题31:讲讲 Flink checkpoint 原理?对齐式和非对齐式 checkpoint 有什么区别?
Checkpoint 作用?
- 保证 Flink 集群在某个算子因为某些原因(如异常退出)出现故障时,能够将整个应用流图的状态恢复到故障之前的某一状态,保证应用流图状态的一致性。Checkpoint 是一种容错恢复机制。
Checkpoint 保存的是什么数据?
- 当前检查点开始时数据源(例如 Kafka)中消息的 offset。
- 记录了所有有状态的 operator 当前的状态信息(例如 sum 中的数值)。
Checkpoint 有两种实现方式:对齐式(Aligned Checkpoint)和非对齐式(Unaligned Checkpoint)。
-
对齐式 Checkpoint:
- 计算所有执行中的任务完成当前状态后,最终整个程序的一个完整状态。
- 取得一个全局会话锁,暂停所有输入数据源的操作,等待所有任务的结果输出。
- 对任务进行 Barrier 插入,通过 Barrier Barrier 来将任务切分成 Snapshotable 和 Non-Snapshotable 两类任务。Snapshotable 任务需要将其状态发送到其他 TaskManager 进行二次备份,而 Non-Snapshotable 任务则不需要。
- 在所有任务都完成 Snapshotable 操作之后,JobManager 根据接收到的各个任务的实际状态,重新计算出恢复点位置。
- 恢复以该状态为准的下一个 CheckPoint,之前的 CheckPoint 使用完毕并且作废。
-
非对齐式 Checkpoint:
- 每个任务在被触发 Checkpoint 时,都记录下自己当前的状态信息。
- 在所有任务完成状态保存之后,JobManager 会选择其中任意一个 Checkpoint 作为重启点。当它恢复时,每个任务将自己记录的状态发送给它所属的 Operator 进行恢复。
- 非对齐式 Checkpoint 允许各个 Task 在不同的时间点异步进行 Checkpoint 操作,节省了执行任务的总体时间。
总结:
- 对齐式 Checkpoint 可以保证所有任务的状态是一致的,但是需要等待所有任务都完成 Checkpoint 后才能进入下一个 Checkpoint,因此会影响整个应用程序的处理速度,exactly once 精确一次性支持。
- 而非对齐式 Checkpoint 则可以保证任务相互之间的状态是独立的,每个任务在自己的频率上异步进行 Checkpoint,可以大大提高系统的可扩展性和容错性,支持 at least once 语义最少一次,消息不会丢失,但是可能会重复。
在 Flink 中,Checkpoint 除了帮助我们实现分布式快照外,还可以通过控制 Checkpoint 的间隔、最大并发数、存储位置等参数来优化应用程序的性能,并确保应用程序的容错性。
面试题32:讲讲 watermark 工作机制?
Watermark 的意义:
- 标识 Flink 任务的事件时间进度,从而能够推动事件时间窗口的触发、计算。
- 解决事件时间窗口的乱序问题。
Watermark 的触发时机:
watermark 时间 >= window_end_time即max(timestamp, currentMaxTimestamp....) - allowedLateness >= window_end_time。- 在
[window_start_time, window_end_time)中有数据存在。
乱序处理可归纳为:
- 窗口 window 的作用是为了周期性的获取数据。
- watermark 的作用是防止数据出现乱序(经常),事件时间内获取不到指定的全部数据,而做的一种保险方法。
- allowLateNess 是将窗口关闭时间再延迟一段时间。
- sideOutPut 是最后兜底操作,所有过期延迟数据,指定窗口已经彻底关闭了,就会把数据放到侧输出流。
面试题33:Flink 双流 join?
- Flink window join:
join():- 通俗理解,将两条实时流中元素分配到同一个时间窗口中完成 Join。两条实时流数据缓存在 Window State 中,当窗口触发计算时,执行 join 操作。(窗口对齐才会触发)。
- 支持 Tumbling Window Join(滚动窗口),Sliding Window Join(滑动窗口),Session Widnow Join(会话窗口),支持处理时间和事件时间两种时间特征。
- 源码核心总结:windows 窗口 + state 存储 + 双层 for 循环执行 join()。
coGroup():- coGroup 的作用和 join 基本相同,但有一点不一样的是,如果未能找到新到来的数据与另一个流在 window 中存在的匹配数据,仍会将其输出。
- 只有 inner join 肯定还不够,如何实现 left/right outer join 呢?答案就是利用 coGroup() 算子。
- 它的调用方式类似于 join() 算子,也需要开窗,但是 CoGroupFunction 比 JoinFunction 更加灵活,可以按照用户指定的逻辑匹配左流和/或右流的数据并输出。(二重循环
end
