使用PySpark在HDFS上重新分区数据失败问题求助
问题分析与解决方案
一、代码操作的核心问题
你同时使用repartition("date", "hour")和partitionBy("date", "hour")属于冗余操作,且直接引发了后续的执行故障:
repartition("date", "hour")会将数据按这两个字段分成若干Spark分区,后续partitionBy写入时,每个Spark分区的数据会直接写入对应(date,hour)的子目录,这会导致每个子目录下只有1个文件。如果某个时间分区的数据量极大,单个Executor需要处理远超负载的数据集,极易引发内存溢出或IO瓶颈,进而导致Executor进程崩溃。- 这种写法完全没达到“让文件大小均衡”的目的,反而可能制造出超大文件。
二、环境配置的潜在问题
- HDFS副本数限制:默认HDFS副本数为3,而你的环境刚好只有3个DataNode。一旦某个DataNode因负载过高宕机,HDFS找不到额外的备用节点来完成副本写入,就会抛出
Failed to replace a bad datanode...的错误。 - Spark资源分配不合理:每个DataNode是8核24GB内存,若Spark默认的Executor资源配置(内存/核数)不合理,会导致节点资源被占满,引发进程崩溃或IO阻塞。
三、具体修复方案
1. 调整Spark代码,实现文件大小均衡
方案一:指定总分区数,结合partitionBy
根据总数据量74GB,目标每个Parquet文件控制在128MB-256MB左右,计算总分区数约290-580个,示例代码:
( df .repartition(300) # 根据目标文件大小调整数值 .write .partitionBy("date", "hour") .mode("overwrite") .format("parquet") .save(OUTPUT_DIRECTORY) )
这样每个(date,hour)子目录下会有多个大小均衡的文件。
方案二:用参数控制单文件大小
如果想更精准控制单文件大小,可使用maxRecordsPerFile参数(按记录数限制):
( df .write .partitionBy("date", "hour") .option("maxRecordsPerFile", 500000) # 根据单条记录大小调整数值 .mode("overwrite") .format("parquet") .save(OUTPUT_DIRECTORY) )
2. 调整HDFS配置,避免副本节点不足
临时降低副本数到2(若业务允许),这样即使一个DataNode故障,仍有备用节点可用:
hdfs dfs -setrep -R 2 hdfs://192.168.1.60:9000/my_repartitioned_dataset
若要长期生效,修改hdfs-site.xml中的dfs.replication为2,重启HDFS服务。
3. 优化Spark资源配置
提交任务时指定合理的Executor资源,避免节点过载:
spark-submit --master spark://<data-node-ip>:7077 \ --executor-cores 4 \ --executor-memory 8G \ --total-executor-cores 24 \ your_script.py
每个DataNode运行2个Executor(4核8G),既充分利用资源,又不会占满节点。
4. 额外排查点
- 查看DataNode日志(默认路径
$HADOOP_HOME/logs/hadoop-*-datanode-*.log),确认是否因磁盘IO过高、内存不足导致节点宕机; - 检查DataNode的磁盘剩余空间,避免因磁盘满引发写入失败。
内容的提问来源于stack exchange,提问作者CopyOfA
相关产品推荐
相关产品推荐

