请阐述 Spark Streaming 中数据分区的基本原理以及并行处理的具体实现方式。
考察说明
考查对 Spark Streaming 数据分区与并行处理机制的理解。
回答思路
- 【回答框架 1】Spark Streaming 的数据分区核心是 DStream,其底层由一系列连续的 RDD 构成,每个 RDD 又划分为多个分区,分区是并行处理的基本单位。分区数决定了任务并行度,通常与输入源的并行度相关,如 Kafka 分区数。
- 【回答框架 2】并行处理通过微批次实现,每个批次的时间间隔内的数据被封装为一个 RDD,该 RDD 的分区会被分配给不同的 Executor 上的任务并行计算。Spark 的任务调度器会根据分区数启动相应数量的 Task,每个 Task 处理一个分区的数据。
- 【回答框架 3】调整并行度的方法包括设置 spark.default.parallelism 参数、使用 repartition 或 coalesce 算子显式调整分区数,以及优化输入源的并行度。合理设置分区数需权衡资源利用与调度开销。
- 【回答框架 4】分区数过少会导致资源闲置、处理延迟;过多则增加任务调度和网络传输开销。实践中需根据数据量、集群资源和延迟要求动态调整。
- 【关键点 1】DStream 的每个 RDD 分区数决定并行度。
- 【关键点 2】并行度可通过 repartition、coalesce 及输入源分区数调整。
- 【关键点 3】分区数设置需平衡资源利用与调度开销。
- 【易错点 1】不要将分区数与并行度混淆,后者还受 Executor 核数限制。
- 【易错点 2】Kafka 分区数不是唯一决定因素,还需考虑接收批次大小。