如何使用PySpark检查Blob存储中Excel文件的必填列是否存在
PySpark实现Blob存储Excel文件必填列校验
前置依赖
确保SparkSession初始化时包含以下依赖包:
- Excel读取依赖:
com.crealytics:spark-excel_2.12:0.13.7(版本需匹配你的Spark版本) - Azure Blob存储依赖:
org.apache.hadoop:hadoop-azure:3.3.4
实现步骤与代码
- 初始化SparkSession并配置Blob存储访问
替换代码中的<your-storage-account>和<your-storage-key>为实际的Azure存储账户名和密钥。
from pyspark.sql import SparkSession from pyspark.sql.utils import AnalysisException # 初始化SparkSession spark = SparkSession.builder \ .appName("ExcelColumnValidation") \ .config("spark.jars.packages", "com.crealytics:spark-excel_2.12:0.13.7,org.apache.hadoop:hadoop-azure:3.3.4") \ .config("fs.azure.account.key.<your-storage-account>.blob.core.windows.net", "<your-storage-key>") \ .getOrCreate()
- 定义各文件的必填列映射
以字典形式存储Blob中每个Excel文件的路径和对应的必填列列表:
required_columns_map = { "wasbs://<container-name>@<your-storage-account>.blob.core.windows.net/excel1.xlsx": ["a", "b", "c"], "wasbs://<container-name>@<your-storage-account>.blob.core.windows.net/excel2.xlsx": ["d", "e", "f"] }
- 遍历文件并校验列
仅读取Excel表头(不加载全量数据),检查必填列是否存在,缺失则抛出断言错误:
for file_path, required_cols in required_columns_map.items(): try: # 仅读取一行数据获取列名,减少资源占用 df = spark.read \ .format("com.crealytics.spark.excel") \ .option("header", "true") \ .option("inferSchema", "false") \ .option("maxRowsInMemory", 1) \ .load(file_path) # 统一列名大小写,避免大小写不匹配问题 actual_columns = [col.lower() for col in df.columns] required_columns_lower = [col.lower() for col in required_cols] # 筛选缺失的必填列 missing_columns = [col for col in required_columns_lower if col not in actual_columns] if missing_columns: raise AssertionError(f"文件 {file_path} 缺少必填列:{', '.join(missing_columns)}") print(f"文件 {file_path} 列校验通过") except AnalysisException as e: print(f"读取文件 {file_path} 失败:{str(e)}") except AssertionError as e: # 抛出断言错误终止程序,或根据业务需求调整处理逻辑 raise e
关键说明
maxRowsInMemory=1:仅加载一行数据,大幅降低内存消耗,因为我们只需要列名信息- 列名大小写统一:避免Excel中列名大小写与定义的必填列不一致导致的误判
- Blob路径格式:遵循
wasbs://<容器名>@<存储账户名>.blob.core.windows.net/<文件路径>的格式
内容的提问来源于stack exchange,提问作者Swati B
相关产品推荐
相关产品推荐

