Spark读取Parquet小分区文件时生成大量任务的优化问询
问题原因与解决方案
为什么生成近20000个任务
Spark默认根据输入文件的分片数生成任务,每个输入分片对应一个任务。你的场景里:
- 按
date_hour分区,一年的时间范围对应约8760个分区,每个分区有15-18个30-70KB的小文件,总文件数接近20000; - 小文件远小于Spark默认的
spark.sql.files.maxPartitionBytes(128MB),每个小文件会被当作独立的输入分片,因此任务数等于符合条件的文件总数。
如何让单个任务读取整个分区内容提升效率
可以通过以下几种方式优化,减少任务数并提升效率:
1. 写入阶段合并小文件(推荐)
从根源解决问题,在写入分区数据时,强制每个分区生成少量大文件:
- 使用
repartition按分区字段重分区后再写入:
这样每个分区只会生成与Spark并行度匹配的文件数(默认是集群核数),避免大量小文件。df.repartition(col("date_hour")) .write .partitionBy("date_hour") .mode("overwrite") .parquet("s3://your-bucket/path") - 或者使用
coalesce减少分区数后写入,适合数据量不大的场景:
每个分区会生成1个大文件。df.coalesce(1) .write .partitionBy("date_hour") .mode("overwrite") .parquet("s3://your-bucket/path")
2. 读取阶段合并小文件
如果无法重新写入数据,可通过调整Spark参数让读取时合并多个小文件为一个输入分片:
- 调大
spark.sql.files.maxPartitionBytes:设置一个较大的值(比如100MB),让Spark将多个小文件合并成一个分片,减少任务数:spark.conf.set("spark.sql.files.maxPartitionBytes", "104857600") # 100MB - 调大
spark.sql.files.openCostInBytes:这个参数用于衡量打开文件的开销,调大后Spark会更倾向于合并小文件(认为打开多个小文件的开销更高):spark.conf.set("spark.sql.files.openCostInBytes", "16777216") # 16MB
3. 利用Parquet元数据加速count统计
你的查询是按分区统计count,可开启Parquet元数据缓存和统计信息利用,避免读取文件内容:
- 开启元数据缓存:
spark.conf.set("spark.sql.parquet.cacheMetadata", "true") - 确保Parquet过滤下推和向量化读取开启(Spark2.4.7默认已开启):
开启后Spark可以直接从Parquet文件的元数据中获取记录数,无需扫描整个文件,大幅提升count查询效率。spark.conf.set("spark.sql.parquet.filterPushdown", "true") spark.conf.set("spark.sql.parquet.enableVectorizedReader", "true")
内容的提问来源于stack exchange,提问作者Amil
相关产品推荐
相关产品推荐

