PySpark加载XML文件时如何指定嵌套Struct中PostalCode列为StringType
方法一:自定义完整嵌套结构Schema
PySpark的StructType支持任意层级的嵌套定义,你只需要按照现有嵌套结构逐层编写Schema,将PostalCode指定为StringType即可,完整示例如下:
from pyspark.sql.types import StructType, StructField, StringType, DateType, DoubleType # 定义Address嵌套结构体Schema address_schema = StructType([ StructField("City", StringType(), True), StructField("Country", StringType(), True), StructField("Number", StringType(), True), StructField("PostalCode", StringType(), True), # 显式指定为字符串类型,保留前导零 StructField("StreetName", StringType(), True) ]) # 定义BankAccount嵌套结构体Schema bank_account_schema = StructType([ StructField("BankAccountNumber", StringType(), True), StructField("CurrencyCode", StringType(), True) ]) # 定义Company结构体Schema,嵌套上述两个子结构体 company_schema = StructType([ StructField("Address", address_schema, True), StructField("BankAccount", bank_account_schema, True) ]) # 定义最外层完整表Schema full_schema = StructType([ StructField("AuditFileCountry", StringType(), True), StructField("AuditFileDateCreated", DateType(), True), StructField("AuditFileVersion", DoubleType(), True), StructField("Company", company_schema, True) ])
加载XML时直接传入上述自定义Schema即可:
# 需替换为你实际的XML行级标签、文件路径 df = spark.read.format("xml") \ .option("rowTag", "你的行级根标签") \ .schema(full_schema) \ .load("xml文件路径")
方法二:修改自动推断的Schema(适用于字段多、结构复杂的场景)
如果Schema字段量级很大,全量手写成本高,可以先读取少量样本推断出基础Schema,再递归修改指定嵌套字段的类型:
from pyspark.sql.types import StringType # 加载1%样本快速推断原始Schema infer_sample_df = spark.read.format("xml") \ .option("rowTag", "你的行级根标签") \ .load("xml文件路径", samplingRatio=0.01) original_schema = infer_sample_df.schema # 递归修改嵌套字段类型的工具函数 def modify_nested_field(schema, field_path, new_type): path_segments = field_path.split(".") new_schema = StructType() for field in schema.fields: if field.name == path_segments[0]: if len(path_segments) == 1: # 匹配到目标字段,替换类型 new_schema.add(StructField(field.name, new_type, field.nullable)) else: # 嵌套结构,递归处理下一层级 new_schema.add(StructField( field.name, modify_nested_field(field.dataType, ".".join(path_segments[1:]), new_type), field.nullable )) else: new_schema.add(field) return new_schema # 修改指定嵌套字段的类型为StringType final_schema = modify_nested_field(original_schema, "Company.Address.PostalCode", StringType()) # 用修改后的Schema加载全量数据 final_df = spark.read.format("xml") \ .option("rowTag", "你的行级根标签") \ .schema(final_schema) \ .load("xml文件路径")
内容的提问来源于stack exchange,提问作者Farrukh Tashpulatov
相关产品推荐
相关产品推荐

