读取含多工作表与合并单元格的Excel文件(Python/PySpark)
解决含合并单元格的Excel转PySpark DataFrame的痛点问题
一、优先用Pandas预处理再转PySpark(最高效方案)
PySpark原生读Excel对合并单元格、复杂表头的支持有限,先通过Pandas完成清洗再转PySpark是最省心的路径。
1. 修复合并单元格的空值填充问题
合并单元格在Excel中仅左上角有值,其余位置为空。如果ffill()没生效,大概率是读取时未正确识别空值范围,直接针对目标列强制填充:
import pandas as pd from pyspark.sql import SparkSession # 先不指定表头,读取所有内容 df_pandas = pd.read_excel("your_file.xlsx", header=None) # 针对前5列单独做向下填充,确保合并单元格值覆盖所有行 df_pandas.iloc[:, :5] = df_pandas.iloc[:, :5].ffill(axis=0) # 其他列按需填充 df_pandas = df_pandas.ffill(axis=0)
2. 动态定位表头,跳过顶部无效标题行
不用硬编码skiprows,通过遍历行识别真正的表头(可自定义判断逻辑):
header_idx = None # 遍历每行,找到第一个包含有效列名的行(示例:非空值≥5个) for idx, row in df_pandas.iterrows(): if not row.isnull().all() and len(row.dropna()) >= 5: header_idx = idx break # 设定表头并过滤之前的无效行 df_pandas.columns = df_pandas.iloc[header_idx] df_pandas = df_pandas[header_idx+1:].reset_index(drop=True)
3. 跳过最后一行无效数据
判断最后一行是否为无效行(比如全空、合计行),直接删除:
last_row = df_pandas.iloc[-1] # 可自定义判断条件,比如包含"合计"关键词或全空 if last_row.isnull().all() or "合计" in str(last_row.values): df_pandas = df_pandas.iloc[:-1]
4. 自动生成PySpark Schema,无需手动编写
转PySpark时直接利用自动推断,或基于Pandas dtype生成精准Schema:
spark = SparkSession.builder.appName("ExcelToSpark").getOrCreate() # 方法1:自动推断Schema(快速便捷) df_spark = spark.createDataFrame(df_pandas) # 方法2:基于Pandas dtype生成精准Schema(适合类型要求严格的场景) from pyspark.sql.types import StructType, StructField, StringType, IntegerType, FloatType def pandas_to_spark_dtype(dtype): if pd.api.types.is_integer_dtype(dtype): return IntegerType() elif pd.api.types.is_float_dtype(dtype): return FloatType() else: return StringType() schema = StructType([ StructField(col, pandas_to_spark_dtype(df_pandas[col].dtype), True) for col in df_pandas.columns ]) df_spark = spark.createDataFrame(df_pandas, schema=schema)
二、纯PySpark读取的替代方案(需依赖第三方包)
如果不想用Pandas,可借助crealytics的Spark Excel插件(需先在Spark中引入依赖),但处理合并单元格要借助窗口函数:
# 读取Excel,不指定表头 df_spark = spark.read.format("com.crealytics.spark.excel") \ .option("header", "false") \ .option("inferSchema", "true") \ .load("your_file.xlsx") # 用窗口函数向下填充前5列的合并单元格空值 from pyspark.sql.window import Window from pyspark.sql.functions import last, col window_spec = Window.partitionBy().orderBy(df_spark.columns[0]).rowsBetween(Window.unboundedPreceding, 0) for col_name in df_spark.columns[:5]: df_spark = df_spark.withColumn(col_name, last(col(col_name), ignorenulls=True).over(window_spec)) # 后续表头定位、最后一行删除逻辑与Pandas类似,需先过滤出表头行再重构DataFrame
内容的提问来源于stack exchange,提问作者Pysparker
相关产品推荐
相关产品推荐

