在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的最大记录处理限制:
- 进入Databricks集群配置页面,在Spark > 高级选项 > Spark config中添加:
spark.driver.extraJavaOptions -Dpoi.maxRecords=200000000 spark.executor.extraJavaOptions -Dpoi.maxRecords=200000000
- 重启集群后,重新执行加载代码:
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
相关产品推荐
相关产品推荐

