如何仅用PySpark从Blob存储加载指定前缀的XML文件?
我在Databricks上开发一个从Azure Blob存储加载XML文件的项目,不同前缀的文件对应不同处理逻辑,因此为每种前缀的文件创建了独立Notebook,示例文件如Account_20240101.xml。
此前通过OS库实现了从Blob存储加载指定前缀的文件,但队友认为该方式效率低下,要求改用纯PySpark实现。原OS库代码可正常运行,但纯PySpark代码执行时出现凭据错误,无法访问Blob存储。
报错信息
Py4JJavaError: An error occurred while calling o412.load.
: shaded.databricks.org.apache.hadoop.fs.azure.AzureException: shaded.databricks.org.apache.hadoop.fs.azure.AzureException: Unable to access container incoming in account *******.blob.core.windows.net using anonymous credentials, and no credentials found for them in the configuration
核心问题是Spark会话的Blob存储凭据配置错误,同时文件路径未匹配原代码中的XML_Test/目录,以下是修正后的实现方式:
1. 正确配置Blob存储凭据
原代码中Spark会话配置使用了占位符<your-storage-account-name>,未替换为实际存储账户名,导致凭据无法生效。在Databricks中,推荐使用spark.conf.set动态设置存储访问密钥,配置项格式需严格对应存储账户域名。
2. 修正文件路径
原OS库代码从XML_Test/{prefix}前缀下读取文件,但纯PySpark代码遗漏了XML_Test/目录,导致路径指向错误,需在Blob路径中补充该目录。
3. 文件过滤优化
若需排除包含特定子串的文件,可在加载后通过input_file_name()函数过滤,若文件数量较多,也可通过更精确的Glob模式在加载前过滤,减少不必要的文件读取。
修正后的纯PySpark代码
from pyspark.sql import SparkSession from pyspark.sql.functions import regexp_extract, input_file_name # 初始化Spark会话 spark = SparkSession.builder \ .appName("Load Multiple XML Files to DataFrame") \ .getOrCreate() # Azure Blob Storage 配置信息 storage_account_name = "******" storage_account_key = "*******" container_name = "incoming" prefix = "accounts" excluded_substring = "cards-snapshots" # 设置Blob存储访问密钥 spark.conf.set( f"fs.azure.account.key.{storage_account_name}.blob.core.windows.net", storage_account_key ) # 构建正确的Blob存储路径,包含XML_Test目录与前缀匹配规则 blob_storage_path = f"wasbs://{container_name}@{storage_account_name}.blob.core.windows.net/XML_Test/{prefix}*" # 加载XML文件、过滤排除项、添加文件名标识列 raw_df = spark.read.format("com.databricks.spark.xml") \ .option("rowTag", "AccountSnapshot") \ .load(blob_storage_path) \ .filter(~input_file_name().contains(excluded_substring)) \ .withColumn( "@FileName", regexp_extract(input_file_name(), r"([^/]+)$", 1) # 提取文件名 ) # 展示数据样本 display(raw_df.limit(10))
额外优化建议
- 安全凭据管理:若Databricks集群已通过集群初始化脚本、Azure AD服务主体或SAS令牌配置了Blob存储访问权限,可移除代码中硬编码的存储密钥,提升安全性。
- Schema预定义:针对大量XML文件,提前定义Schema并通过
.schema()指定,避免Spark多次扫描文件推断Schema,大幅提升加载效率:from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 自定义Schema示例 account_schema = StructType([ StructField("AccountID", StringType(), nullable=True), StructField("AccountName", StringType(), nullable=True), StructField("Balance", IntegerType(), nullable=True) # 根据实际XML结构补充其他字段 ]) # 加载时指定Schema raw_df = spark.read.format("com.databricks.spark.xml") \ .option("rowTag", "AccountSnapshot") \ .schema(account_schema) \ .load(blob_storage_path)
内容的提问来源于stack exchange,提问作者dexon

