PySpark DataFrame写入HDFS Parquet时部分分区文件缺失致丢数排查
PySpark写入Parquet时分区文件数不足且记录丢失的排查与解决
可能原因及对应排查/解决手段
静默数据兼容性问题
Parquet对数据类型有严格要求,若DataFrame中存在不兼容类型(比如复杂嵌套结构、超出范围的数值),Spark可能会静默丢弃部分数据而非报错。- 开启严格校验:设置
spark.sql.strictChecks=true和spark.sql.parquet.writeLegacyFormat=false,强制触发不兼容场景的报错,而非静默丢弃。 - 检查任务级日志:查看每个Executor的stdout/stderr日志,可能存在数据截断、类型转换失败的警告,这类信息不会体现在DAG的任务状态中,但会导致数据丢失。
- 开启严格校验:设置
HDFS文件系统一致性延迟
任务执行完成不代表HDFS上的文件已完全落地或同步,可能存在副本未完成写入就被标记为成功的情况。- 写入后执行HDFS完整性检查:运行
hdfs fsck /your/target/path -files -blocks,确认所有文件的块状态正常。 - 避免并发写入冲突:确保写入路径没有其他进程同时进行写入/删除操作,覆盖写入时明确指定
mode="overwrite"。
- 写入后执行HDFS完整性检查:运行
Spark自适应执行的分区合并
若开启了自适应执行(spark.sql.adaptive.enabled=true),Spark可能会根据数据量自动合并分区,导致最终文件数少于初始分区数。- 禁用自适应分区合并:设置
spark.sql.adaptive.coalescePartitions.enabled=false,或直接关闭自适应执行spark.sql.adaptive.enabled=false,保持初始分区数不变。
- 禁用自适应分区合并:设置
分区目录异常
部分分区目录可能因权限问题、HDFS异常被意外移动或删除,DAG仅标记任务执行完成,不校验最终文件存在性。- 写入后校验文件数量:遍历目标路径下的Parquet文件,统计数量是否为100,定位缺失的分区目录,排查对应目录的HDFS权限或操作日志。
快速验证方法
写入完成后立即对比原DataFrame和写入后数据的记录数,确认是否真的丢失:
original_count = df.count() df.write.parquet("/your/target/path") written_count = spark.read.parquet("/your/target/path").count() print(f"原记录数: {original_count}, 写入后记录数: {written_count}")
若数量不一致,重点排查记录数差值对应的任务日志,定位具体丢失的分区数据。
内容的提问来源于stack exchange,提问作者Prateek Pathak
相关产品推荐
相关产品推荐

