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

在PySpark中展平嵌套XML及关联数据的技术咨询

嵌套XML转平展DataFrame及关联方案

1. 展平嵌套DataFrame的函数与示例

针对你的嵌套结构,主要用到PySpark的explode(展开数组类型字段)、col(提取嵌套结构体字段)、select(选择与重命名字段)函数。以下是完整实现步骤:

步骤说明

  1. 读取XML文件(需依赖spark-xml库)
  2. 展开数组类型字段:先展开个人名称数组NAME_LIST.PEC_INDVDL_NAME,再展开注册状态数组ENRLMTS.ABC_855X.ENRLMT_INFO.ENRLMT_DTLS.ENRLMT_STUS_DLTS
  3. 逐层提取嵌套结构体中的所有字段,重命名为清晰的列名
  4. 过滤或整理冗余字段(可选)

代码示例

from pyspark.sql import SparkSession
from pyspark.sql.functions import explode, col

# 初始化SparkSession
spark = SparkSession.builder \
    .appName("FlattenXML") \
    .config("spark.jars.packages", "com.databricks:spark-xml_2.12:0.16.0") \
    .getOrCreate()

# 读取XML文件,指定根节点为PRVDR
df = spark.read \
    .format("xml") \
    .option("rootTag", "PRVDR") \
    .option("rowTag", "PRVDR") \
    .load("path/to/your/xml/file.xml")

# 第一步:展开个人名称数组
df_with_names = df.select(
    col("PRVDR_INFO.INDVDL_INFO.*"),
    col("ENRLMTS.*"),
    explode(col("PRVDR_INFO.INDVDL_INFO.NAME_LIST.PEC_INDVDL_NAME")).alias("name_details")
).drop("NAME_LIST")

# 第二步:展开注册状态数组
df_with_status = df_with_names.select(
    col("*"),
    explode(col("ABC_855X.ENRLMT_INFO.ENRLMT_DTLS.ENRLMT_STUS_DLTS")).alias("enrlmt_status")
).drop("ENRLMT_STUS_DLTS")

# 第三步:提取所有嵌套字段并整理列名
flattened_df = df_with_status.select(
    # 个人基本信息
    col("CNTL_ID"),
    col("BIRTH_DT"),
    col("BIRTH_STATE_CD"),
    col("BIRTH_STATE_NAME"),
    col("BIRTH_CNTRY_CD"),
    col("BIRTH_CNTRY_NAME"),
    col("BIRTH_FRGN_SW"),
    # 个人名称信息
    col("name_details.NAME_CD").alias("NAME_CD"),
    col("name_details.NAME_DESC").alias("NAME_DESC"),
    col("name_details.FIRST_NAME").alias("FIRST_NAME"),
    col("name_details.MDL_NAME").alias("MDL_NAME"),
    col("name_details.LAST_NAME").alias("LAST_NAME"),
    col("name_details.TRMNTN_DT").alias("TRMNTN_DT"),
    col("name_details.DATA_STUS_CD").alias("NAME_DATA_STUS_CD"),
    # TIN信息
    col("PEC_TIN.TIN"),
    col("PEC_TIN.TAX_IDENT_TYPE_CD"),
    col("PEC_TIN.TAX_IDENT_DESC"),
    col("PEC_TIN.DATA_STUS_CD").alias("TIN_DATA_STUS_CD"),
    # NPI信息
    col("PEC_NPI.NPI"),
    col("PEC_NPI.VRFYD_BUSNS_SW"),
    col("PEC_NPI.CREAT_TS"),
    col("PEC_NPI.DATA_STUS_CD").alias("NPI_DATA_STUS_CD"),
    # 注册基本信息
    col("ABC_855X.ACPT_NEW_PTNT_SW"),
    col("ABC_855X.ENRLMT_INFO.ENRLMT_DTLS.FORM_TYPE_CD"),
    col("ABC_855X.ENRLMT_INFO.ENRLMT_DTLS.ENRLMT_ID"),
    col("ABC_855X.ENRLMT_INFO.ENRLMT_DTLS.BUSNS_STATE"),
    col("ABC_855X.ENRLMT_INFO.ENRLMT_DTLS.BUSNS_STATE_NAME"),
    # 合同方信息
    col("ABC_855X.ENRLMT_INFO.ENRLMT_DTLS.CNTRCTR_LIST.CNTRCTR_INFO.CNTRCTR_ID"),
    col("ABC_855X.ENRLMT_INFO.ENRLMT_DTLS.CNTRCTR_LIST.CNTRCTR_INFO.CNTRCTR_NAME"),
    col("ABC_855X.ENRLMT_INFO.ENRLMT_DTLS.CNTRCTR_LIST.CNTRCTR_INFO.DATA_STUS_CD").alias("CNTRCTR_DATA_STUS_CD"),
    # 注册状态信息
    col("enrlmt_status.STUS_CD").alias("ENRLMT_STUS_CD"),
    col("enrlmt_status.STUS_DESC").alias("ENRLMT_STUS_DESC"),
    col("enrlmt_status.STUS_DT").alias("ENRLMT_STUS_DT"),
    col("enrlmt_status.DATA_STUS_CD").alias("ENRLMT_STUS_DATA_STUS_CD"),
    # 注册状态原因信息
    col("enrlmt_status.ENRLMT_STUS_RSN_DLTS.STUS_RSN_CD").alias("STUS_RSN_CD"),
    col("enrlmt_status.ENRLMT_STUS_RSN_DLTS.STUS_RSN_DESC").alias("STUS_RSN_DESC"),
    col("enrlmt_status.ENRLMT_STUS_RSN_DLTS.DATA_STUS_CD").alias("STUS_RSN_DATA_STUS_CD"),
    # 重新验证信息
    col("ABC_855X.PEC_ENRLMT_REVLDTN.REVLDTN_INSTNC_NUM"),
    col("ABC_855X.PEC_ENRLMT_REVLDTN.REVLDTN_STUS_CD"),
    col("ABC_855X.PEC_ENRLMT_REVLDTN.REVLDTN_STUS_DESC")
)

# 展示平展后的结果
flattened_df.show(truncate=False)

2. 关联PRVDR_INFO与ENRLMTS记录

方法一:使用CNTL_ID直接关联

CNTL_ID是PRVDR_INFO中的主键,每个PRVDR节点对应一组ENRLMTS数据。在上述平展过程中,CNTL_ID会被保留到每一行记录中,自然实现PRVDR_INFO与ENRLMTS的关联——所有从同一PRVDR节点展开的行,都会带有相同的CNTL_ID,无需额外关联操作。

如果需要将PRVDR_INFO和ENRLMTS拆分为两个独立的表再关联:

# 提取PRVDR_INFO表(去重)
prvdr_info_df = flattened_df.select(
    "CNTL_ID", "BIRTH_DT", "BIRTH_STATE_CD", "BIRTH_STATE_NAME",
    "BIRTH_CNTRY_CD", "BIRTH_CNTRY_NAME", "BIRTH_FRGN_SW",
    "TIN", "TAX_IDENT_TYPE_CD", "TAX_IDENT_DESC", "TIN_DATA_STUS_CD",
    "NPI", "VRFYD_BUSNS_SW", "CREAT_TS", "NPI_DATA_STUS_CD"
).distinct()

# 提取ENRLMTS表(保留CNTL_ID)
enrlmts_df = flattened_df.select(
    "CNTL_ID",
    "ACPT_NEW_PTNT_SW", "FORM_TYPE_CD", "ENRLMT_ID",
    "BUSNS_STATE", "BUSNS_STATE_NAME", "CNTRCTR_ID", "CNTRCTR_NAME",
    "CNTRCTR_DATA_STUS_CD", "ENRLMT_STUS_CD", "ENRLMT_STUS_DESC",
    "ENRLMT_STUS_DT", "ENRLMT_STUS_DATA_STUS_CD", "STUS_RSN_CD",
    "STUS_RSN_DESC", "STUS_RSN_DATA_STUS_CD", "REVLDTN_INSTNC_NUM",
    "REVLDTN_STUS_CD", "REVLDTN_STUS_DESC"
)

# 关联两个表
joined_df = prvdr_info_df.join(enrlmts_df, on="CNTL_ID", how="inner")

方法二:动态生成代理键(当无主键时)

如果CNTL_ID不可用或不存在,可以用monotonically_increasing_id()生成全局唯一的代理键,为每个原始PRVDR行分配键值,再展开ENRLMTS时保留该键:

from pyspark.sql.functions import monotonically_increasing_id

# 生成代理键
df_with_proxy = df.withColumn("PROXY_KEY", monotonically_increasing_id())

# 后续展开步骤同前,保留PROXY_KEY字段即可实现关联
# ...(省略展开代码,只需在select中包含PROXY_KEY)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 20:20:24