配置spark.sql.files.maxPartitionBytes后分区倾斜的解决方法咨询
问题场景
在pyspark Docker容器中操作,尝试将本地Parquet文件读取到Spark并拆分多个分区,操作步骤如下:
- 下载大小为46Mb的NY Taxi数据集:
curl -Lo taxi.parquet 'https://d37ci6vzurychx.cloudfront.net/trip-data/yellow_tripdata_2023-01.parquet'
- 启动Spark会话,配置最大分区大小为10Mb:
spark = ( SparkSession .builder .appName("Bucketing") .master("local[4]") .config("spark.sql.files.maxPartitionBytes", 10 * 1000 * 1024) .getOrCreate() )
- 将Parquet文件读取为DataFrame:
taxi_df = spark.read.parquet("taxi.parquet")
- 触发操作:
taxi_df.show()
执行后发现:Spark UI显示生成了4个分区,但所有数据都集中在一个分区,其余三个为空,出现分区倾斜。根据Spark官方文档,spark.sql.files.maxPartitionBytes应在读取文件时生效,但实际未按预期工作。额外测试了36Mb和113Mb的Parquet文件,分区数计算正确(分别为3和13),但数据仍全部集中在单个分区,排除了local[4]配置的影响。
核心原因
Parquet是列式存储格式,单个Parquet文件的内部按行组(Row Group)划分,Spark读取时会将每个行组作为一个基本分区单元,无法拆分单个行组。即使设置了spark.sql.files.maxPartitionBytes,如果文件的所有数据都在一个行组里,Spark也只能将整个文件作为一个分区加载——这就是分区倾斜的根源,你测试的这些NY Taxi Parquet文件本身都是单一行组结构。
解决方法
1. 读取后手动重分区
如果需要强制拆分分区,可在读取DataFrame后调用repartition()(coalesce仅能减少分区数,增加分区需用repartition):
# 按指定数量重分区 taxi_df = taxi_df.repartition(4) # 或按某列分区(适合后续需按该列操作的场景) taxi_df = taxi_df.repartition("passenger_count")
该方式会触发Shuffle,适合数据量不大的场景。
2. 预处理Parquet文件,拆分多行组
从根源解决的话,可通过Spark重新写入文件时配置更小的行组大小(例如10Mb),让文件包含多个行组,这样Spark读取时会自动拆分成分区:
# 写入时设置行组大小为10Mb taxi_df.write.option("parquet.block.size", 10 * 1024 * 1024).parquet("split_taxi.parquet") # 读取拆分后的文件 split_taxi_df = spark.read.parquet("split_taxi.parquet")
3. 调整Spark分区配置(仅适用于多行组文件)
如果文件本身有多行组,但spark.sql.files.maxPartitionBytes未生效,可调整spark.sql.files.openCostInBytes(默认4Mb),调小该值会让Spark更倾向于拆分分区:
spark = ( SparkSession .builder .appName("Bucketing") .master("local[4]") .config("spark.sql.files.maxPartitionBytes", 10 * 1024 * 1024) .config("spark.sql.files.openCostInBytes", 1 * 1024 * 1024) .getOrCreate() )
注意:此配置仅对包含多个行组的Parquet文件有效,单一行组文件依然无法拆分。
内容的提问来源于stack exchange,提问作者neshkeev

