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

在Databricks中用PySpark从ADLS挂载点加载大型Excel遇阻求方案

在Databricks中使用PySpark加载大型Excel文件的替代方案

针对你遇到的pyspark.pandas Arrow转换错误、spark-excel加载大文件时的POI内存限制问题,以下是几种可行的解决方法:

方法1:Pandas UDF分布式分块读取

利用Pandas的分块读取能力结合PySpark Pandas UDF,实现分布式加载,既规避Arrow转换异常,又解决单节点内存瓶颈:

from pyspark.sql.functions import pandas_udf, col
import pandas as pd

# 提前定义目标Schema,示例格式根据你的实际数据调整
SCHEMA = "id int, user_name string, order_amount double, order_date timestamp"

@pandas_udf(SCHEMA)
def read_excel_chunk(file_path_series: pd.Series) -> pd.DataFrame:
    file_path = file_path_series.iloc[0]
    # 按指定块大小分读取Excel,chunksize可根据节点内存调整
    chunks = pd.read_excel(file_path, engine="openpyxl", chunksize=100000)
    return pd.concat(chunks)

# 生成包含文件路径的RDD,触发分布式读取逻辑
file_rdd = spark.sparkContext.parallelize(["dbfs:/mnt/aadata/ds/data/test.xlsx"])
df = file_rdd.toDF(["file_path"]).select(read_excel_chunk(col("file_path")))

# 验证加载结果
df.printSchema()
df.show(5)

方法2:转换为Parquet/CSV后加载(生产环境推荐)

Excel并非为大数据场景设计,转换为列存格式(如Parquet)可大幅提升Spark读取效率与稳定性,以下是两种转换方式:

驱动节点分块转换(适用于文件大小在驱动内存范围内)

import pandas as pd

# 分块读取Excel并逐块写入Parquet
chunks = pd.read_excel("dbfs:/mnt/aadata/ds/data/test.xlsx", engine="openpyxl", chunksize=100000)
for idx, chunk in enumerate(chunks):
    chunk.to_parquet(f"dbfs:/mnt/aadata/ds/data/test_chunk_{idx}.parquet")

# 合并所有Parquet文件为Spark DataFrame
df = spark.read.parquet("dbfs:/mnt/aadata/ds/data/test_chunk_*.parquet")

分布式批量转换(超大型文件)

若文件过大无法在驱动节点处理,可先通过外部工具(如Azure Data Factory)将大Excel拆分为多个小文件,再用spark-excel批量读取后转存:

# 读取拆分后的小Excel文件
df = spark.read.format("com.crealytics.spark.excel") \
        .option("header", "true") \
        .option("schema", SCHEMA) \
        .load('dbfs:/mnt/aadata/ds/data/split_excels/*.xlsx')

# 写入Parquet格式
df.write.mode("overwrite").parquet("dbfs:/mnt/aadata/ds/data/test_parquet")

方法3:调整spark-excel的POI参数解决内存限制

针对org.apache.poi.util.RecordFormatException错误,可通过修改Spark集群JVM参数提升POI的最大记录处理限制:

  1. 进入Databricks集群配置页面,在Spark > 高级选项 > Spark config中添加:
spark.driver.extraJavaOptions -Dpoi.maxRecords=200000000
spark.executor.extraJavaOptions -Dpoi.maxRecords=200000000
  1. 重启集群后,重新执行加载代码:
df=spark.read.format("com.crealytics.spark.excel") \
        .option("header", "true") \
        .option("schema", SCHEMA) \
        .load('dbfs:/mnt/aadata/ds/data/test.xlsx')

注:poi.maxRecords值需根据文件实际大小调整,避免超出节点内存承载范围。

内容的提问来源于stack exchange,提问作者Lambo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 04:35:47