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

如何合并同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),避免单节点内存过载):
    df.coalesce(1).write.mode("overwrite").parquet("Downloads/new/export")
    
    Spark采用懒加载机制,仅在执行写入操作时才处理数据,不会将全量数据加载到内存,效率远高于Pandas。

二、解决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环境变量

  1. 下载与Spark版本匹配的winutils工具包(例如Spark 3.x对应Hadoop 3.3.x版本)
  2. 解压后设置系统环境变量:
    • HADOOP_HOME:指向解压后的winutils目录
    • PATH:追加%HADOOP_HOME%\bin
  3. 重启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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 09:30:49