PySpark写入DataFrame为CSV时仅Driver工作Executor不运行问题
问题诱发原因
- 配置未生效:代码中单独定义的
SparkConf对象没有传入SparkSession构建流程,且未指定集群部署模式、Executor资源参数,默认启动的是本地单线程模式,所有计算逻辑都运行在Driver节点,不会向集群申请Executor资源。 - 数据分片不足:如果HDFS上存储的Parquet文件本身只有少量文件块(极端情况为单个大文件),Spark读取后生成的DataFrame分区数和文件块数一致,没有足够的任务分片分发到Executor,即使Executor正常启动也无法分配到计算任务。
- 写入逻辑未优化:读取后未对DataFrame做重分区调整,写入CSV时并行度和读取时一致,低并行度下只会占用少量计算资源,无法发挥集群分布式计算能力。
修复方案
- 修正SparkSession初始化逻辑,正确配置集群参数,确保Executor能正常注册并参与计算,修正后的初始化代码如下:
import findspark findspark.init() from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("parquet_to_csv") \ # 根据集群类型选择部署模式,standalone/yarn/k8s,禁止使用默认local模式 .master("yarn") \ # 按集群实际资源配置Executor规格与数量 .config("spark.executor.cores", "4") \ .config("spark.executor.memory", "8g") \ .config("spark.executor.instances", "10") \ # 开启动态分区调整,自动优化任务并行度 .config("spark.sql.adaptive.enabled", "true") \ .config("spark.sql.adaptive.coalescePartitions.enabled", "true") \ .getOrCreate()
- 读取Parquet文件后手动调整分区数,保证任务分片数匹配集群总计算能力:
df = spark.read.format("parquet").load("hdfs://hadoop1:9001/data.parquet") # 分区数设置为集群总CPU核心数的2~3倍即可,例:10个Executor每个4核,总核数40,分区数设为80~120 df = df.repartition(100)
- 正常执行写入逻辑即可,禁止在写入前调用
coalesce(1)强行合并为单分区,否则会将所有计算拉到单个节点执行,退化为单机计算:
df.write.format("csv").mode("overwrite").save("hdfs://hadoop1:9001/data.csv")
额外排查项
- 任务提交必须使用
spark-submit命令执行脚本,直接本地运行Python脚本默认不会连接集群,所有计算都在本地执行。 - 任务启动后访问Spark UI的Executors页面,确认所有配置的Executor都正常注册,若Executor注册失败先排查集群资源剩余量、节点间通信、权限配置问题。
- 若HDFS上的Parquet为单个远大于块大小的文件,建议先对Parquet做重分区存储,避免读取时无法拆分分片导致并行度不足。
- 如果最终需要输出单个CSV文件,不要在Spark中合并分区,待多分区CSV写入完成后,使用HDFS自带的
hdfs dfs -getmerge命令合并文件即可,避免单点计算压力。
内容的提问来源于stack exchange,提问作者刘一汉
相关产品推荐
相关产品推荐

