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

PySpark/Pandas批量读取多Excel文件 合并同名单表生成Delta表

批量处理多Excel文件并生成Delta表方案

实现步骤与代码

1. 依赖准备

确保已安装必要依赖库:

pip install pandas openpyxl pyspark delta-spark

2. 核心逻辑实现

import pandas as pd
from glob import glob
from pyspark.sql import SparkSession
from delta.tables import DeltaTable

# 初始化Spark Session(用于Delta表操作)
spark = SparkSession.builder \
    .appName("ExcelToDelta") \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
    .getOrCreate()

# 初始化存储容器:key为工作表名,value为对应DataFrame列表
merged_dfs = {}
# 存储每个文件的reference数据(用于和对应data表关联)
file_ref_map = {}

# 遍历所有Excel文件
excel_paths = glob("/path/to/your/excel/files/*.xlsx")  # 替换为实际文件路径
for path in excel_paths:
    # 读取当前文件所有工作表
    all_sheets = pd.read_excel(path, sheet_name=None)
    
    # 读取固定的reference表
    ref_df = all_sheets["reference"]
    file_ref_map[path] = ref_df
    
    # 筛选data类工作表(匹配以data开头的表名,忽略大小写)
    data_sheets = [sheet for sheet in all_sheets.keys() if sheet.lower().startswith("data")]
    
    # 处理每个data工作表
    for sheet_name in data_sheets:
        data_df = all_sheets[sheet_name]
        # 与reference表关联(替换为你的实际关联逻辑)
        processed_df = pd.merge(data_df, ref_df, how="left", on="关联列名")  # 替换为实际关联键
        
        # 将处理后的DataFrame加入对应列表
        if sheet_name not in merged_dfs:
            merged_dfs[sheet_name] = []
        merged_dfs[sheet_name].append(processed_df)

# 合并同名称工作表
final_merged = {}
for sheet_name, df_list in merged_dfs.items():
    final_merged[sheet_name] = pd.concat(df_list, ignore_index=True)
    print(f"已合并{len(df_list)}个文件的{sheet_name}表,共{len(final_merged[sheet_name])}条数据")

# 写入Delta表
delta_base_path = "/path/to/delta/tables"  # 替换为Delta表存储路径
for sheet_name, df in final_merged.items():
    # 将pandas DataFrame转为Spark DataFrame
    spark_df = spark.createDataFrame(df)
    # 定义Delta表名称(如data1对应Customers_Data1)
    delta_table_name = f"Customers_{sheet_name.capitalize()}"
    delta_table_path = f"{delta_base_path}/{delta_table_name}"
    
    # 写入Delta表(若表已存在则追加,否则创建)
    spark_df.write \
        .format("delta") \
        .mode("append") \
        .save(delta_table_path)
    
    print(f"已成功生成Delta表:{delta_table_name}")

关键说明

  • 工作表筛选:用sheet.lower().startswith("data")兼容不同大小写的表名(如Data1、data_2等),可根据实际规则调整筛选条件
  • 关联逻辑:示例用pd.merge实现data表与reference表的关联,需替换为你的实际业务关联规则(如自定义函数)
  • Delta表写入:采用append模式支持增量写入,若需覆盖现有表可改为mode("overwrite")
  • 批量文件读取:用glob批量获取所有Excel文件路径,支持通配符匹配

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 11:05:20