如何合并同Schema的Parquet文件?Spark写入报错求助
合并同Schema Parquet文件及Spark写入错误解决
一、Spark合并Parquet的高效方法
完全不用把数据加载到Pandas内存里,Spark本身支持分布式读取合并,且能完美保留原Schema,适配后续创建Delta表的需求:
- 批量读取文件:若所有Parquet文件在同一目录,直接读取目录即可,比逐个指定文件更便捷:
df = spark.read.parquet("Downloads/parquet_files/") - 指定目标文件:如果文件分散,可使用通配符或手动传入文件路径列表:
# 通配符匹配所有以P开头的Parquet文件 df = spark.read.parquet("Downloads/P*.parquet") # 手动传入文件路径列表 file_paths = ["Downloads/P4.parquet", "Downloads/P3.parquet", "Downloads/P2.parquet", "Downloads/P1.parquet"] df = spark.read.parquet(*file_paths) - 写入单个Parquet文件:用
coalesce(1)将数据合并为单个文件(适合小数据量;大数据量建议用repartition(1),避免单节点内存过载):
Spark采用懒加载机制,仅在执行写入操作时才处理数据,不会将全量数据加载到内存,效率远高于Pandas。df.coalesce(1).write.mode("overwrite").parquet("Downloads/new/export")
二、解决Windows下的UnsatisfiedLinkError错误
你遇到的是Windows环境中Hadoop原生依赖缺失导致的问题,三种快速修复方案:
方案1:禁用Hadoop原生IO
在创建SparkSession时添加配置,跳过原生库检查:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("MergeParquet") \ .config("spark.hadoop.io.native.lib.available", "false") \ .getOrCreate()
方案2:配置Hadoop环境变量
- 下载与Spark版本匹配的winutils工具包(例如Spark 3.x对应Hadoop 3.3.x版本)
- 解压后设置系统环境变量:
HADOOP_HOME:指向解压后的winutils目录PATH:追加%HADOOP_HOME%\bin
- 重启Notebook或Python环境后重新执行代码
方案3:直接写入Delta表(跳过中间Parquet步骤)
既然最终目标是创建Delta表,可直接读取Parquet后写入Delta,省去合并Parquet的环节,同时绕开该错误:
df = spark.read.parquet("Downloads/P*.parquet") # 直接生成Delta表 df.write.mode("overwrite").format("delta").save("Downloads/new/delta_table")
后续可直接基于该Delta表进行处理,流程更高效。
内容的提问来源于stack exchange,提问作者martin
相关产品推荐
相关产品推荐

