改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

50道大数据精选面试题

ckckck2025年1月10日71 浏览

Hive 面试题

面试题1:Hive 中四个 by 的区别?

  1. Sort By:分区内有序;

    • 不是全局排序,其在数据进入 reducer 前完成排序,也就是说它会在数据进入 reduce 之前为每个 reducer 都产生一个排序后的文件。因此,如果用 sort by 进行排序,并且设置 mapreduce.job.reduces > 1,则 sort by 只保证每个 reducer 的输出有序,不保证全局有序。
  2. Order By:全局排序,只有一个 Reducer;

    • order by 会对输入做全局排序,因此只有一个 reducer(多个 reducer 无法保证全局有序),然而只有一个 reducer,会导致当输入规模较大时,消耗较长的计算时间。
  3. Distribute By:类似 MR 中 Partition,进行分区,结合 sort by 使用。

    • distribute by 是控制在 map 端如何拆分数据给 reduce 端的。类似于 MapReduce 中分区 partitioner 对数据进行分区。
    • Hive 会根据 distribute by 后面列,将数据分发给对应的 reducer,默认是采用 hash 算法 + 取余数的方式。
  4. 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)

    1. 静态分区是在编译期间指定的指定分区名。
    2. 支持 load 和 insert 两种插入方式。
      • load 方式
        1. 会将分区字段的值全部修改为指定的内容。
        2. 一般是确定该分区内容是一致的时候才会使用。
      • insert 方式
        1. 必须先将数据放在一个没有设置分区的普通表中。
        2. 该方式可以在一个分区内存储一个范围的内容。
        3. 从普通表中选出的字段不能包含分区字段。
    3. 适用于分区数少,分区名可以明确的数据。
  • 动态分区 DP(dynamic partition)

    1. 根据分区字段的实际值,动态进行分区。
    2. 是在 SQL 执行的时候进行分区。
    3. 需要先将动态分区设置打开(set hive.exec.dynamic.partition.mode=nonstrict)。
    4. 只能用 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 端:

  1. 关闭自动 offset,手动提交 offset。
    • 设置 enable.auto.commit = false,默认值 true,自动提交。
    • 使用 Kafka 的 Consumer 的类,用方法 consumer.commitSync() 提交。
    • 或者使用 spring-kafka 的 Acknowledgment 类,用方法 ack.acknowledge() 提交(推荐使用)。
  2. 另一个方法同样需要手动 commit offset,另外在 consumer 端再将所有 fetch 到的数据缓存到 queue 里,当把 queue 里所有的数据处理完之后,再批量提交 offset,这样就能保证只有处理完的数据才被 commit。

面试题5:Kafka 如何保证数据 exactly-once?

  1. Producer exactly-once

    • enable.idempotence=true
    • 分区副本数 >= 2
    • isr >=2
    • ProducerID + SequenceNumber + Ack = -1(幂等性)。
  2. Consumer exactly-once

    • 手动维护并提交偏移量。
    • 设置 enable.auto.commit=false,关闭自动提交偏移量。
    • 借助外部数据库,如 Redis 的 pipeline,MySQL 的事务机制管理存储偏移量。
    • 在同一事物中,在消息被处理完之后再提交偏移量。并更新偏移量。否则消息需回滚,并获取到上一次偏移量的位置重新进行处理。

面试题6:Kafka 数据积压怎么解决?

  1. 增加 broker 节点,增加分区数量,提高并行度。
  2. 修改单个消费为批量消费。
  3. 增加单线程消费为线程池异步消费。
  4. 缩短批次时间间隔。
  5. 老版本 SparkStreaming 控制消费的速率(spark.streaming.kafka.maxRatePerPartition),可以控制最大的消费速率,在参数中设置;新版本设置背压机制实现消费处理的动态平衡。
  6. 对代码进行优化,尽可能的一次性计算多个结果,减少 shuffle 过程。
  7. 处理的结果如果过多,可以将数据保存到 MySQL 集群、MongoDB 集群【支持事物】或 ES【不支持事物】,增大吞吐量。
  8. 消费线程将拉取的消息放到一个滑动窗口中,通过滑动窗口控制拉取的速度。
  9. 对于倾斜的 key 加以处理,加随机数等方式打散。

面试题7:Kafka rebalance 发生时机和分区分配策略?

Rebalance 触发时机?

  • 当出现以下几种情况时,Kafka 会进行一次重新分区分配操作,即 Kafka 消费者端的 Rebalance 操作:
    1. 同一个 consumer 消费者组 group.id 中,新增了消费者进来,会执行 Rebalance 操作。
    2. 消费者离开当期所属的 consumer group 组。比如主动停机或者宕机。
    3. 分区数量发生变化时(即 topic 的分区数量发生变化时)。
    4. 消费者主动取消订阅。

Rebalance 三种策略?

  • Kafka 新版本提供了三种 rebalance 分区分配策略:(partition.assignment.strategy):
    1. range
    2. round-robin
    3. 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 通常与应用消费消息后进行什么处理的关系更大,相对复杂一些。

分区数过多的危害?

  1. 客户端/服务器端需要使用的内存就越多。
  2. 文件句柄的开销。
  3. 越多的分区可能增加端对端的延迟。
  4. 降低高可用性。

面试题9:数据发往 Kafka 的分区规则?

key 和 value 的类型,一般都用字符串即可。数据到底写入到哪一个分区中:

  • 如果指定了分区,就写入到指定的分区中。
  • 如果没有指定分区,指定了 key,按照 key 的 hashcode,取模,写入对应的分区。
  • 如果没有指定分区和 key,轮询机制。

面试题10:Kafka producer buffer pool 的作用?

Kafka 通过使用内存缓冲池的设计,让整个发送过程中的存储空间循环利用,有效减少 JVM GC 造成的影响,从而提高发送性能,提升吞吐量。

面试题11:Kafka 时间轮的作用?

Kafka 通过时间轮来处理延迟任务,只将时间轮的槽保存到延迟队列,大大的减少了延迟队列的元素数量,这样对于元素的增加删除性能有很大提高;Kafka 通过阻塞的方式 poll 延迟队列的,减少了大量的空转;为了保证线程安全,灵活运用读写锁、原子对象、synchronized 控制时间轮的操作。

面试题12:Kafka 为什么这么快?

Kafka 为什么这么快

Spark 面试问题

面试题13:Spark 为什么比 MapReduce 快?

这是一道常见的面试题,回答时可以从 IO、shuffle 与排序、资源、部署模式、内存管理策略等各个方面来回答。

  1. MR 基于磁盘的分布式计算引擎,频繁的磁盘 IO。Spark 基于内存进行计算,DAG 计算模型,大大减少了磁盘 IO。
  2. Spark 多线程运行,MR 多进程运行。
  3. Spark 粗粒度资源申请,MR 细粒度资源申请。
  4. Spark 支持多种部署模式,MR 只支持 yarn 上部署。
  5. Shuffle 与排序;MR 有 reducer 必排序,一般会经过 3 次排序,Spark Shuffle 数据的排序操作不是必须的。
    • Spark 有多种 shuffle 类型,Spark 不一定会发生 shuffle,MR 一定会发生 shuffle。
  6. Spark 具有灵活的内存管理策略。

面试题14:Spark Repartition 和 Coalesce 的关系与区别,能简单说说吗?

  1. 关系

    • 两者都是用来改变 RDD 的 partition 数量的,repartition 底层调用的就是 coalesce 方法:coalesce(numPartitions, shuffle = true)
  2. 区别

    • repartition 一定会发生 shuffle,coalesce 根据传入的参数来判断是否发生 shuffle。
    • 一般情况下增大 RDD 的 partition 数量使用 repartition,减少 partition 数量时使用 coalesce。

面试题15:简述下 Spark 中的缓存(cache 和 persist)与 checkpoint 机制,并指出两者的区别和联系?

关于 Spark 缓存和检查点的区别,大致可以从这 4 个角度去回答:

  1. 位置
    • Persist 和 Cache 将数据保存在内存,Checkpoint 将数据保存在 HDFS。
  2. 生命周期
    • Persist 和 Cache 程序结束后会被清除或手动调用 unpersist 方法,Checkpoint 永久存储不会被删除。
  3. RDD 依赖关系
    • Persist 和 Cache,不会丢掉 RDD 间的依赖链/依赖关系,CheckPoint 会斩断依赖链。
  4. 执行与使用
    • persist 中 RDD 的逻辑只会执行一次,而 checkpoint 会执行两次。
    • 生产环境中一般都是 cache 和 checkpoint 连用,这样 RDD 逻辑只会执行一次,并且会缓存到 checkpoint 中。
Spark 缓存与 checkpoint

面试题16:Spark on Yarn client 模式与 cluster 的区别?

  1. driver 所在位置不同
    • client 模式下 driver 线程只在 spark-submit 命令提交的机器上。
    • cluster 模式下,driver 线程只在 applicationMaster 所在的节点。
  2. 启动的任务进程名字不一样
    • client 模式下:ExecutorLauncher 只负责向 yarn 申请容器来启动 executor。
    • cluster 模式下,applicationMaster 既要负责申请运行 executor 的资源,又要调 Driver 线程来做 task 调度。

面试题17:Spark RDD、DataFrame、Dataset 的区别与联系?

三者的共性:

  1. RDD、DataFrame、DataSet 全都是 Spark 平台下的分布式弹性数据集,为处理超大型数据提供便利。
  2. 三者都有惰性机制,在进行创建、转换,如 map 方法时,不会立即执行,只有在遇到 Action 如 foreach 时,三者才会开始遍历运算。
  3. 三者有许多共同的函数,如 filter,排序等。
  4. 在对 DataFrame 和 Dataset 进行操作许多操作都需要这个包:import spark.implicits._(在创建好 SparkSession 对象后尽量直接导入)。
  5. 三者都会根据 Spark 的内存情况自动缓存运算,这样即使数据量很大,也不用担心会内存溢出。
  6. 三者都有 partition 的概念。
  7. DataFrame 和 Dataset 均可使用模式匹配获取各个字段的值和类型。

三者的区别:

  1. RDD
    • RDD 一般和 Spark MLlib 同时使用。
    • RDD 不支持 SparkSQL 操作。
  2. DataFrame
    • 与 RDD 和 DataSet 不同,DataFrame 每一行的类型固定为 Row,每一列的值没法直接访问,只有通过解析才能获取各个字段的值。
    • DataFrame 与 DataSet 一般不与 Spark MLlib 同时使用。
    • DataFrame 与 DataSet 均支持 SparkSQL 的操作,比如 select,groupby 之类,还能注册临时表/视窗,进行 SQL 语句操作。
    • DataFrame 与 DataSet 支持一些特别方便的保存方式,比如保存成 csv,可以带上表头,这样每一列的字段名一目了然。
  3. DataSet
    • DataSet 和 DataFrame 拥有完全相同的成员函数,区别只是每一行的数据类型不同。DataFrame 其实就是 DataSet 的一个特例。type DataFrame = DataSet[Row]
    • DataFrame 也可以叫 DataSet[Row],每一行类型是 Row,不解析,每一行究竟有哪些字段,各个字段又是什么类型都无从得知,只能用上面的 getAs 方法或者共性中的第七条提到的模式匹配拿出特定字段,而 DataSet 中,每一行是什么类型是不一定的,在自定义 case class 之后可以很自由的获取每一行的信息。

三者的转换:

  1. 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()
    • 方案三:构造 Schema 的方式:.createDataFrame(rdd, scheme)
  2. DataFrame 转 RDD
    • df.rdd 或者 df.javaRDD()
  3. RDD 转 DataSet
    • 方案一:使用 toDS() 算子,需要导入隐式转换(import spark.implicits._)。
    • 方案二:使用 spark.createDataset(rdd)
  4. DataSet 转 RDD
    • 直接使用 .rdd
  5. DataFrame 转 DataSet
    • 封装样例类,调用 df.as[xxx]
    • case class xxx()
    • df.as[xxx]
  6. DataSet 转 DataFrame
    • ds.toDF()
Spark RDD、DataFrame、Dataset 的区别与联系

面试题18:updateStateByKey 与 mapWithState 使用区别?

  • updateStateByKey:统计全局的 key 的状态,就算没有数据输入,它也会在每一个批次的时候返回之前的 key 的状态。
    • 缺点:若数据量太大的话,需要 checkpoint 的数据会占用较大的存储,效率低下。
  • mapWithState:也是用于全局统计 key 的状态,但是它如果没有数据输入,便不会返回之前的 key 的状态,有一点增量的感觉。效率更高,生产中建议使用。
    • 优点:我们可以只是关心那些已经发生的变化的 key,对于没有数据输入,则不会返回那些没有变化的 key 的数据。这样的话,即使数据量很大,checkpoint 也不会像 updateStateByKey 那样,占用太多的存储。

面试题19:Spark SQL 三种 join 方式?

  1. Broadcast Hash Join:适合一张很小的表和一张大表进行 Join。
  2. Shuffle Hash Join:适合一张小表(比上一个大一点)和一张大表进行 Join。
  3. Sort Merge Join:适合两张大表进行 Join。

Shuffle Hash Join 策略必须满足以下条件:

  1. 仅支持等值 Join,不要求参与 Join 的 Keys 可排序(这点是和 sort-merge join 相对应)。
  2. spark.sql.join.preferSortMergeJoin 参数必须设置为 false,参数是从 Spark 2.0.0 版本引入的,默认值为 true,也就是默认情况下选择 Sort Merge Join。
  3. 小表的大小(plan.stats.sizeInBytes)必须小于 spark.sql.autoBroadcastJoinThreshold * spark.sql.shuffle.partitions(默认值 200)其实就是让每一个小表的分区都类似于广播变量的小表。
  4. 而且小表大小(stats.sizeInBytes)的三倍必须小于等于大表的大小(stats.sizeInBytes),也就是 a.stats.sizeInBytes * 3 <= b.stats.sizeInBytes

Broadcast Hash Join 策略必须满足以下条件:

  1. 小表的数据必须很小,可以通过 spark.sql.autoBroadcastJoinThreshold 参数来配置,默认是 10MB。
  2. 如果内存比较大,可以将阈值适当加大。
  3. spark.sql.autoBroadcastJoinThreshold 参数设置为 -1,可以关闭这种连接方式。
  4. 只能用于等值 Join,不要求参与 Join 的 keys 可排序。

要启用 Shuffle Sort Merge Join 必须满足的条件是仅支持等值 Join,并且要求参与 Join 的 Keys 可排序。

面试题20:RDD 有什么缺陷?

  1. 不支持细粒度的写和更新操作(如网络爬虫),Spark 写数据是粗粒度的。所谓粗粒度,就是批量写入数据,为了提高效率。但是读数据是细粒度的,也就是说可以一条条的读。
  2. 不支持增量迭代计算,Flink 支持。

面试题21:groupByKey 和 reduceByKey 区别?

reduceByKey 和 groupByKey 都存在 shuffle 的操作,但是 reduceByKey 可以在 shuffle 前对分区内相同 key 的数据进行预聚合(combine)功能,这样会减少落盘的数据量,而 groupByKey 只是进行分组,不存在数据量减少的问题,reduceByKey 性能比较高。从功能的角度:reduceByKey 其实包含分组和聚合的功能。GroupByKey 只能分组,不能聚合,所以在分组聚合的场合下,推荐使用 reduceByKey,如果仅仅是分组而不需要聚合。那么还是只能使用 groupByKey。

面试题22:RDD 的弹性表现在哪几点?

  1. 自动的进行内存和磁盘的存储切换。
  2. 基于 Lineage 的高效容错。
  3. Task 如果失败会自动进行特定次数的重试。
  4. Stage 如果失败会自动进行特定次数的重试,而且只会计算失败的分片。
  5. Checkpoint 和 persist,数据计算之后持久化缓存。
  6. 数据调度弹性,DAG task 调度和资源无关。
  7. 数据分片的高度弹性。

面试题23:RDD 通过 Lineage(记录数据更新)的方式为何很高效?

  1. Lazy 记录了数据的来源,RDD 是不可变的,且是 lazy 级别的,且 RDD 之间构成了链条,lazy 是弹性的基石。由于 RDD 不可变,所以每次操作就产生新的 RDD,不存在全局修改的问题,控制难度下降,所有有计算链条将复杂计算链条存储下来,计算的时候从后往前回溯 900 步是上一个 stage 的结束,要么就 checkpoint。
  2. 记录原数据,是每次修改都记录,代价很大如果修改一个集合,代价就很小,官方说 RDD 是粗粒度的操作,是为了效率,为了简化,每次都是操作数据集合,写或者修改操作,都是基于集合的 RDD 的写操作是粗粒度的,RDD 的读操作既可以是粗粒度的也可以是细粒度,读可以读其中的一条条的记录。
  3. 简化复杂度,是高效率的一方面,写的粗粒度限制了使用场景如网络爬虫,现实世界中,大多数写是粗粒度的场景。

面试题24: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 的两种机制:
    1. 估算,去要内存 5.01 * 2 - 5,要不到的时候就去排序,最终溢写的小的磁盘小文件合并成为了一个大的磁盘小文件。
    2. 当不需要排序的时候,默认使用 Bypass 机制。
  • Bypass 运行机制的触发条件:
    • Shuffle reduce task 数量小于 spark.shuffle.sort.bypassMergeThreshold 参数的值小于 200,不开启,溢写磁盘不需要排序,小于等于的时候是开启的。
    • 不是聚合类的 shuffle 算子(比如 reduceByKey)。

总结:

  • Hash shuffle(合并运行机制)优化机制产生的磁盘小文件的个数:C * R(core * reducer)。
  • Hash shuffle(普通):产生的磁盘小文件:M * R
  • Sort shuffle 产生的磁盘小文件的个数为:2 * M
  • Bypass 机制产生的磁盘小文件的个数为:2 * M

面试题26:哪些 Spark 算子会有 shuffle 过程?

  • 去重:distinct
  • 排序:groupByKeyreduceByKeysortByKey
  • 重分区:repartitionrepartitionAndSortWithinPartitionscoalesce
  • 集合或者表连接操作:joincogroup

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 错误。禁用算子链可以避免此类问题。
  • 需要流数据处理与批数据处理共存:如果同时需要进行流式与批处理,禁用算子链可以让流和批处理同时运行,避免出现串行化的问题。

Keyby 之后出现数据倾斜常见原因?

  • Key 的选择不合适:如果选择的 Key 不平衡或者有明显的热点数据,就容易出现数据倾斜的问题。应该尝试选择更加平衡的 Key,例如多个属性组合的方式。
  • 数据分布不均匀:有些数据在时间、空间上分布不均匀,导致某些 Key 的数据量比其他 Key 大很多。可以通过统计每个 Key 对应的数据量,找到数据分布不均匀的原因。
  • 算子链长/复杂度高:当算子链过长或者算子的操作很复杂时,也容易导致某些 Task 的数据处理量过大。可以通过拆分算子链、优化算子操作等方式来解决。
  • 并行度设置不当:并行度过高可能导致资源浪费,过低则会导致数据倾斜。应该根据实际情况,合理设置并行度。

定位:

  1. 步骤 1:定位反压
    • 定位反压有 2 种方式:Flink Web UI 自带的反压监控(直接方式)、Flink Task Metrics(间接方式)。通过监控反压的信息,可以获取到数据处理瓶颈的 Subtask。
  2. 步骤 2:确定数据倾斜
    • Flink Web UI 自带 Subtask 接收和发送的数据量。当 Subtasks 之间处理的数据量有较大的差距,则该 Subtask 出现数据倾斜。如下图所示,红框内的 Subtask 出现数据热点。

解决方案:

  • keyBy 后聚合操作存在数据倾斜(通过 Flink LocalKeyBy 思想来解决):
    • 在 keyBy 上游算子数据发送之前,首先在上游算子的本地对数据进行聚合后再发送到下游,使下游接收到的数据量大大减少,从而使得 keyBy 之后的聚合操作不再是任务的瓶颈。类似 MapReduce 中 Combiner 的思想,但是这要求聚合操作必须是多条数据或者一批数据才能聚合,单条数据没有办法通过聚合来减少数据量。从 Flink LocalKeyBy 实现原理来讲,必然会存在一个积攒批次的过程,在上游算子中必须攒够一定的数据量,对这些数据聚合后再发送到下游。
    • 注意:Flink 是实时流处理,如果 keyby 之后的聚合操作存在数据倾斜,且没有开窗口的情况下,简单的认为使用两阶段聚合,是不能解决问题的。因为这个时候 Flink 是来一条处理一条,且向下游发送一条结果,对于原来 keyby 的维度(第二阶段聚合)来讲,数据量并没有减少,且结果重复计算(非 FlinkSQL,未使用回撤流)。
  • keyBy 后窗口聚合操作存在数据倾斜(两阶段聚合):
    • 因为使用了窗口,变成了有界数据的处理,窗口默认是触发时才会输出一条结果发往下游,所以可以使用两阶段聚合的方式:
      1. 第一阶段聚合:key 拼接随机数前缀或后缀,进行 keyby、开窗、聚合。
      2. 第二阶段聚合:去掉随机数前缀或后缀,按照原来的 key 及 windowEnd 作 keyby、聚合。
Flink 数据倾斜解决方案

TaskManager、slot、并行度之间的关系:

  • 在 Yarn 集群中 Job 分离模式下,TaskManager 的数量 = ceil(slot 数量 / 并行度)slotNumber >= taskManager * 并行度

TaskManager/slots 与 CPU 的关系:

  • 经验上讲 Slot 的数量与 CPU-core 的数量一致为好。但考虑到超线程,可以让 slotNumber = 2 * cpuCore

slot 与并行度:

  • 一般我们设置 task 的并行度不能超过 slot 的数量。
  • 一个 Task 的并行度等于分配给它的 Slot 个数(前提槽资源充足)。
  • application:每个 job 独享一个集群,job 退出则集群退出。main 方法在集群上运行。
  • session:多个 job 共享集群资源,job 退出集群也不会退出。main 方法在客户端运行。
  • pre-job:每个 job 独享一个集群,job 退出则集群退出。main 方法在客户端运行。

适用场景:

  • Session 模式:一般用来部署那些对延迟非常敏感但运行时长较短的作业,需要频繁提交小 job 的场景。
  • Per-Job 模式:一般用来部署那些长时间运行的作业。
  • Application 模式:综合了两种模式的所有优点,建议生产上适用。
Flink 模式区别

Checkpoint 作用?

  • 保证 Flink 集群在某个算子因为某些原因(如异常退出)出现故障时,能够将整个应用流图的状态恢复到故障之前的某一状态,保证应用流图状态的一致性。Checkpoint 是一种容错恢复机制。

Checkpoint 保存的是什么数据?

  • 当前检查点开始时数据源(例如 Kafka)中消息的 offset。
  • 记录了所有有状态的 operator 当前的状态信息(例如 sum 中的数值)。

Checkpoint 有两种实现方式:对齐式(Aligned Checkpoint)和非对齐式(Unaligned Checkpoint)。

  • 对齐式 Checkpoint

    1. 计算所有执行中的任务完成当前状态后,最终整个程序的一个完整状态。
    2. 取得一个全局会话锁,暂停所有输入数据源的操作,等待所有任务的结果输出。
    3. 对任务进行 Barrier 插入,通过 Barrier Barrier 来将任务切分成 Snapshotable 和 Non-Snapshotable 两类任务。Snapshotable 任务需要将其状态发送到其他 TaskManager 进行二次备份,而 Non-Snapshotable 任务则不需要。
    4. 在所有任务都完成 Snapshotable 操作之后,JobManager 根据接收到的各个任务的实际状态,重新计算出恢复点位置。
    5. 恢复以该状态为准的下一个 CheckPoint,之前的 CheckPoint 使用完毕并且作废。
  • 非对齐式 Checkpoint

    1. 每个任务在被触发 Checkpoint 时,都记录下自己当前的状态信息。
    2. 在所有任务完成状态保存之后,JobManager 会选择其中任意一个 Checkpoint 作为重启点。当它恢复时,每个任务将自己记录的状态发送给它所属的 Operator 进行恢复。
    3. 非对齐式 Checkpoint 允许各个 Task 在不同的时间点异步进行 Checkpoint 操作,节省了执行任务的总体时间。

总结:

  • 对齐式 Checkpoint 可以保证所有任务的状态是一致的,但是需要等待所有任务都完成 Checkpoint 后才能进入下一个 Checkpoint,因此会影响整个应用程序的处理速度,exactly once 精确一次性支持。
  • 而非对齐式 Checkpoint 则可以保证任务相互之间的状态是独立的,每个任务在自己的频率上异步进行 Checkpoint,可以大大提高系统的可扩展性和容错性,支持 at least once 语义最少一次,消息不会丢失,但是可能会重复。

在 Flink 中,Checkpoint 除了帮助我们实现分布式快照外,还可以通过控制 Checkpoint 的间隔、最大并发数、存储位置等参数来优化应用程序的性能,并确保应用程序的容错性。

面试题32:讲讲 watermark 工作机制?

Watermark 的意义:

  • 标识 Flink 任务的事件时间进度,从而能够推动事件时间窗口的触发、计算。
  • 解决事件时间窗口的乱序问题。

Watermark 的触发时机:

  1. watermark 时间 >= window_end_timemax(timestamp, currentMaxTimestamp....) - allowedLateness >= window_end_time
  2. [window_start_time, window_end_time) 中有数据存在。

乱序处理可归纳为:

  • 窗口 window 的作用是为了周期性的获取数据。
  • watermark 的作用是防止数据出现乱序(经常),事件时间内获取不到指定的全部数据,而做的一种保险方法。
  • allowLateNess 是将窗口关闭时间再延迟一段时间。
  • sideOutPut 是最后兜底操作,所有过期延迟数据,指定窗口已经彻底关闭了,就会把数据放到侧输出流。
  1. 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