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

基于PySpark与AWS Glue拆分嵌套XML为关联的双DataFrame

使用PySpark + AWS Glue拆分嵌套XML并建立关联

核心方案

从嵌套DataFrame中提取PRVDR_INFO和ENROLMENT数据,生成代理键作为关联字段,分别存入独立DataFrame,最终通过代理键实现两表关联。

假设嵌套DataFrame结构

假设你转换后的Spark DataFrame(Glue DynamicFrame同理)结构如下:

root
 |-- PRVDR_INFO: struct (nullable = true)
 |    |-- PRVDR_ID: string (nullable = true)
 |    |-- PRVDR_NAME: string (nullable = true)
 |-- ENROLMENT: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- ENROL_ID: string (nullable = true)
 |    |    |-- ENROL_DATE: string (nullable = true)

具体实现代码

1. 生成代理键并拆分PRVDR_INFO表

from pyspark.sql.functions import monotonically_increasing_id

# 假设原始嵌套DataFrame名为nested_df
# 添加全局唯一代理键,作为两表关联依据
df_with_key = nested_df.withColumn("provider_key", monotonically_increasing_id())

# 提取并展开PRVDR_INFO结构,保留代理键
prvdr_df = df_with_key.select(
    "provider_key",
    "PRVDR_INFO.PRVDR_ID",
    "PRVDR_INFO.PRVDR_NAME"
    # 按需添加PRVDR_INFO下的其他字段
)

2. 拆分ENROLMENT表并关联代理键

若ENROLMENT为数组类型,先展开数组再提取字段:

from pyspark.sql.functions import explode

# 展开ENROLMENT数组,同时保留关联用的provider_key
enrolment_df = df_with_key.select(
    "provider_key",
    explode("ENROLMENT").alias("enrolment")
).select(
    "provider_key",
    "enrolment.ENROL_ID",
    "enrolment.ENROL_DATE"
    # 按需添加ENROLMENT下的其他字段
)

3. AWS Glue DynamicFrame适配处理

如果使用Glue原生DynamicFrame,可直接用Relationalize API自动拆分嵌套结构并生成关联键:

from awsglue.transforms import Relationalize
from awsglue.context import GlueContext

glueContext = GlueContext(spark.sparkContext)

# 假设你的DynamicFrame名为dynamic_frame
# 自动拆分嵌套结构,生成带关联键的多个表
relationalized_dfs = Relationalize.apply(
    frame=dynamic_frame,
    staging_path="s3://your-staging-bucket/path/",
    name="root",
    transformation_ctx="relationalize"
)

# 获取拆分后的两个表(名称会根据结构自动生成,例如root_PRVDR_INFO、root_ENROLMENT)
prvdr_df = relationalized_dfs["root_PRVDR_INFO"]
enrolment_df = relationalized_dfs["root_ENROLMENT"]

验证关联有效性

通过join操作验证两表关联是否正常:

joined_df = prvdr_df.join(enrolment_df, on="provider_key", how="inner")
joined_df.show()

注意:若XML中每个PRVDR_INFO对应多条ENROLMENT记录,explode操作是必要的;若为一对一关系,直接提取字段即可。代理键也可改用UUID生成(from pyspark.sql.functions import uuid),避免单调递增ID的潜在冲突问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 00:25:28