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

Spark写入多PARQUET文件至HDFS并发性能异常问题排查

PySpark写入HDFS性能差异问题分析

问题背景

我有一个PySpark任务,需将约1.5TB数据以PARQUET格式写入HDFS,相关Spark参数如下:

Num of executor: 500
Driver memory: 16G
Driver cores: 4
Executor memory: 16G
Executor cores: 4
spark.executor.memoryOverhead=16384

spark.sql.shuffle.partitions=2000
spark.default.parallelism=2000

写入HDFS的伪代码如下:

if table is already exists:
    df.write \
            .mode('append') \
            .format("parquet") \
            .insertInto("{}.{}".format(database, table))
else:
    df.write \
            .mode('overwrite') \
            .format("parquet") \
            .partitionBy(partition_columns) \
            .saveAsTable("{}.{}".format(database, table))

在相同参数下,进行了两组测试,唯一区别是写入前repartition设置的分区数:

  • 测试1:将最终数据重分区为500个分区:df.repartition(500)
    任务运行更快更顺畅,作业间间隔极小,最长间隔约8分钟,总耗时1小时。
  • 测试2:将最终数据重分区为2000个分区:df.repartition(2000)
    首次与二次作业间隔达51分钟,相邻作业间均存在长间隔,整体耗时为测试1的5倍(5小时)。

疑问

  1. 为何写入更多文件会导致性能差异如此巨大?同一数据集的多PARQUET文件不应并发写入吗?
  2. 测试2中,Spark UI显示作业完成后,仍在向HDFS写入PARQUET文件,直至生成2000个文件。尽管我有2000个executor核心(500个executor×4核),且设置了spark.sql.shuffle.partitions=2000和spark.default.parallelism=2000,Spark似乎并未并发写入HDFS,是否可实现并发写入?
  3. Spark UI显示作业完成是否仅指内存中处理完成,而HDFS写入仍在进行,这是否是间隔产生的原因?

另外,底层Hadoop服务是否可能对Spark的写入进行限流?


问题解答

1. 更多文件导致性能暴跌的原因

不是不能并发写入,但小文件过多会触发HDFS的元数据瓶颈和写入放大问题:

  • HDFS的NameNode需要维护每个文件的元数据(权限、位置、块信息等),2000个文件相比500个,元数据操作量直接翻4倍,NameNode的RPC请求会被打满,出现排队延迟。
  • 每个Parquet文件写入时,都会经历"创建临时文件→写入数据→flush→重命名为正式文件"的流程,小文件多了之后,这些操作的总开销会指数级上升。另外,Parquet的列存格式本身需要一定的文件大小才能发挥优势,太小的文件会导致每个文件的元数据占比过高,读写效率下降。
  • 如果你的任务用了partitionBy,且分区列基数大,2000个分区再加上分区列的拆分,实际生成的文件数可能远超过2000,进一步加剧NameNode的压力。

2. 为何2000核没实现并发写入?

理论上2000个核心可以对应2000个并行写入任务,但实际受限于几个关键点:

  • HDFS的写入并发限制:HDFS本身对同时写入的文件数有隐含限制,NameNode处理元数据请求的能力是有限的,当并发写入请求超过阈值,就会被排队处理,看起来就是没有完全并发。
  • Spark的写入阶段调度:Spark的写入作业不是所有分区同时启动的,而是分批次调度。如果executor资源被之前的计算任务占用未释放,或者集群里有其他任务抢占资源,都会导致写入任务无法同时启动2000个。
  • 数据倾斜:如果repartition之后的数据分布不均匀,部分分区数据量极大,这些慢分区会拖整体进度,看起来像是并发度不够。

想要实现更高的并发写入,可以尝试:

  • 调大HDFS NameNode的dfs.namenode.handler.count参数,提升元数据处理能力。
  • 确保executor在写入阶段有足够的空闲资源,比如减少计算阶段的资源占用,或者调整动态资源分配参数。
  • 避免数据倾斜,确保repartition后数据均匀分布。

3. Spark UI的作业完成与HDFS写入的关系

是的,Spark UI显示作业完成,仅代表内存中的计算和 shuffle 阶段完成,HDFS的写入可能还在后台进行:

  • Spark的写入流程是:每个task将数据写入本地磁盘的临时文件,然后再上传到HDFS,最后重命名为正式文件。这个上传和重命名的过程,是在task完成之后的"commit"阶段进行的,这部分操作不会被计入Spark作业的执行时间,但会占用集群资源。
  • 当作业完成后,Driver还会协调各个executor完成最终的文件提交,这部分操作如果遇到NameNode忙,就会出现长时间的间隔,也就是你看到的作业间等待。

关于Hadoop服务限流的可能性

完全有可能,常见的限流场景包括:

  • HDFS的dfs.datanode.max.transfer.threads参数限制了DataNode同时处理的读写请求数,如果这个值太小,会导致写入请求排队。
  • YARN的资源调度限制,如果集群里有其他任务,YARN可能会限制Spark的写入阶段资源,导致并发度上不去。
  • 部分企业的Hadoop集群会配置QoS(服务质量)策略,对Spark任务的写入带宽或请求数进行限流,保障核心业务的资源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 17:50:29