在PySpark中展平嵌套XML及关联数据的技术咨询
嵌套XML转平展DataFrame及关联方案
1. 展平嵌套DataFrame的函数与示例
针对你的嵌套结构,主要用到PySpark的explode(展开数组类型字段)、col(提取嵌套结构体字段)、select(选择与重命名字段)函数。以下是完整实现步骤:
步骤说明
- 读取XML文件(需依赖
spark-xml库) - 展开数组类型字段:先展开个人名称数组
NAME_LIST.PEC_INDVDL_NAME,再展开注册状态数组ENRLMTS.ABC_855X.ENRLMT_INFO.ENRLMT_DTLS.ENRLMT_STUS_DLTS - 逐层提取嵌套结构体中的所有字段,重命名为清晰的列名
- 过滤或整理冗余字段(可选)
代码示例
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
相关产品推荐
相关产品推荐

