You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark 3.5.0写入小Parquet文件至GCS耗时过长求优化方案

问题:Spark 3.5.0 Parquet读转写至GCS耗时异常过高

我们采用本地部署的Spark,从Google Cloud Storage(GCS)读取Parquet文件到DataFrame,再写入GCS的另一个Parquet文件夹,使用的Python代码如下:

DFS_BLOCKSIZE = 512 * 1024 * 1024

spark = SparkSession.builder \
        .appName("test_app_parquet_load") \
        .config("spark.master", "spark://spark-master-svc:7077") \
        .config("spark.driver.maxResultSize", '1g') \
        .config("spark.driver.memory", '1g') \
        .config("spark.executor.cores",4) \
        .config("spark.sql.shuffle.partitions", 16) \
        .config("spark.sql.files.maxPartitionBytes", DFS_BLOCKSIZE) \
        .getOrCreate()

folder="gs://input_folder/input1/key=20240610"

print(f"reading parquet from {folder}")

start_time1 = time.time()
data_df = spark.read.parquet(folder)
end_time1 = time.time()
print(f"Time duration for reading parquet t1: {end_time1 - start_time1}")
 
start_time2 = time.time()
data_df.write.mode("overwrite").parquet("gs://output_folder/output/key=20240610")
end_time2 = time.time()
print(f"Time duration for writing parquet t3: {end_time2 - start_time2}")

spark.stop()

测试发现Spark 2.4.5与3.5.0的写入耗时差异极大:2.4.5读取约10秒、写入约20秒;3.5.0读取耗时相近,但写入仅3MB的文件需近3分钟。进一步测试显示,若读取CSV文件再写入Parquet格式,3.5.0仅需30-40秒,仅Parquet读转写场景在3.5.0中耗时过长,已尝试关闭自适应查询执行模式但无效,需优化写入耗时的Parquet相关配置建议。


优化配置建议

以下是针对Spark 3.5.0 Parquet写至GCS的针对性配置优化,可逐一测试验证效果:

1. 调整Parquet写入批量与元数据缓存

  • spark.sql.parquet.writeBatchSize: 设置为1024(默认512),增大写入时的批量处理数据量,减少小IO操作频次
  • spark.sql.parquet.cacheMetadata: 设为false,避免写入时重复缓存元数据,小文件场景下元数据缓存开销占比会被放大

2. 禁用向量化读写(适配小文件场景)

  • spark.sql.parquet.enableVectorizedReader: 设为false,小文件读转写时,向量化处理的额外开销占比更高,切换至普通读取路径可降低耗时
  • spark.sql.parquet.enableVectorizedWriter: 设为false,关闭向量化写入,避免向量格式转换带来的不必要开销

3. GCS存储传输优化

  • spark.hadoop.fs.gs.outputstream.upload.chunk.size: 设置为67108864(64MB),增大GCS上传分块大小,减少HTTP请求次数
  • spark.hadoop.fs.gs.upload.buffer.size: 调整为16777216(16MB),优化上传缓冲区大小,匹配GCS传输特性

4. 控制输出文件大小

  • spark.sql.files.maxRecordsPerFile: 设置为100000,强制每个输出文件的记录数上限,避免生成过多极小文件,减少GCS文件操作开销
  • spark.sql.parquet.compression.codec: 确认使用snappy(默认),避免使用高开销的压缩算法(如lz4在部分场景下会增加写入耗时)

5. 关闭不必要的Schema验证

  • spark.sql.parquet.validateSchema: 设为false,读转写场景Schema一致,无需重复验证,节省校验时间

内容的提问来源于stack exchange,提问作者user3830120

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.20 19:59:51