Spark 3.1 Java环境下分区控制与Dataproc资源使用率问题咨询
Spark输出文件数控制与相关配置问题解答
输出文件数量的核心决定逻辑
Spark最终写入存储的文件数量,等于执行写操作的DataFrame/RDD的非空分区数,空分区不会生成对应文件。你遇到的coalesce不生效的问题,核心原因是:
- coalesce属于窄依赖操作,仅能减少分区数,且如果你在
coalesce(4500)之后还执行了join、groupBy、orderBy之类的shuffle操作,分区数会被重置为你配置的spark.sql.shuffle.partitions=5500,导致coalesce操作失效。 - 你第一个数据集生成5499个文件,正好对应5500个shuffle分区里有1个空分区;第二个数据集生成3500个文件,是因为5500个shuffle分区里有2000个空分区,没有数据输出。
人工控制输出文件数的常用方案
完全支持人工控制,常见方案如下:
- 调整分区操作放在所有shuffle操作之后、写操作之前,优先用repartition(可增可减分区,数据分布更均匀),仅在不需要数据重分布时用coalesce:
// 正确示例:所有计算逻辑完成后再调整分区,直接写 dataset.groupBy("type").agg(sum("value").alias("total")) .repartition(4500) // 这里调整到目标分区数 .write.parquet("gs://your-bucket/path"); - 按单文件大小/行数控制,无需调整分区,直接通过写参数设置:
// 每个输出文件最多写10000行,自动拆分 dataset.write.option("maxRecordsPerFile", 10000) .parquet("gs://your-bucket/path"); - 如果按字段分区写入,可配合
repartition("分区字段")减少每个子目录下的小文件数量。
相关默认配置说明
直接影响分区数的核心配置如下:
spark.sql.shuffle.partitions:所有shuffle操作完成后的默认分区数,Spark 3.1默认值为200,你已手动修改为5500。spark.default.parallelism:RDD类无父依赖操作的默认分区数,YARN模式下默认值为所有Executor总核心数,按你的配置总核心为250节点 * 2Executor/节点 * 3核/Executor = 1500。spark.sql.files.maxPartitionBytes:读取文件时单个分区的最大容量,默认值为128MB,决定了读入数据的初始分区数。
另外你提到的CPU利用率仅25%的问题,大概率是YARN节点的可用CPU核数配置限制,需要确认yarn.nodemanager.resource.cpu-vcores配置是否放开到了每个节点8核的匹配值。
Java代码获取配置默认值的方法
通过SparkSession的conf对象即可读取所有配置,包括默认值:
import org.apache.spark.sql.SparkSession; public class SparkConfigDemo { public static void main(String[] args) { SparkSession spark = SparkSession.builder() .appName("GetConfigDemo") .getOrCreate(); // 获取shuffle分区数配置,未手动配置则返回默认值200 String shufflePartitions = spark.conf().get("spark.sql.shuffle.partitions"); // 获取默认并行度配置 String defaultParallelism = spark.conf().get("spark.default.parallelism"); // 获取读文件最大分区字节数配置 String maxPartitionBytes = spark.conf().get("spark.sql.files.maxPartitionBytes"); // 读取可选配置,不存在时返回自定义默认值 String customConfig = spark.conf().getOption("spark.custom.config").orElse("default_val"); } }
内容的提问来源于stack exchange,提问作者Sweety
相关产品推荐
相关产品推荐

