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

基于PySpark递归展开嵌套数据并生成技术ID的方案问询

PySpark嵌套数据自动展开并生成关联技术ID的可扩展解决方案

问题背景与需求

我从源系统接收JSON文件后先存储为嵌套Delta表,再用PySpark展开嵌套数据。目前手动编写展开代码,未采用递归逻辑,但随着业务市场数量突破20个且各市场Schema不同,手动方式已完全不可行。

需要实现满足以下要求的可扩展方案:

  • 基于PySpark开发
  • 将嵌套数组拆分至独立的DataFrame
  • 自动生成可用于关联各展开表的技术ID(当前因未使用该ID已引发数据关联问题)

需求受flatterer包启发,但该包仅支持Python和原始JSON文件,无法直接用于PySpark场景。现寻求帮助:

  1. 是否存在可用于PySpark的、能展开嵌套数据并自动生成关联技术ID的工具包?
  2. 若无相关包,提供PySpark嵌套数据展开的实践代码。

当前手动实现示例(存在无技术ID的问题)

嵌套Schema与测试数据

from pyspark.sql.types import StructType, StructField, StringType, IntegerType, ArrayType
import pyspark.sql.functions as F

# 定义嵌套Schema
schema = StructType([
    StructField("some_business_id", StringType(), True),
    StructField("some_field", StringType(), True),
    StructField("businesspartners", ArrayType(
        StructType([
            StructField("name", StringType(), True),
            StructField("profession", StringType(), True),
            StructField("addresses", ArrayType(
                StructType([
                    StructField("street", StringType(), True),
                    StructField("city", StringType(), True)
                ])
            ), True)
        ])
    ), True)
])

# 测试数据
data = [
    ("1", "value1", [
        {"name": "claus", "profession": "data engineer", "addresses": [{"street": "Main St", "city": "Metropolis"}]},
        {"name": "joern", "profession": "data engineer", "addresses": [{"street": "Side St", "city": "Gotham"}]}
    ]),
    ("2", "value2", [
        {"name": "claus", "profession": "data scientist", "addresses": [{"street": "Sixt St", "city": "Manhatten"}]},
        {"name": "joern", "profession": "Football player", "addresses": [{"street": "Baker St", "city": "Fullham"}]}
    ]),
]

# 创建初始嵌套DataFrame
df = spark.createDataFrame(data, schema=schema)

手动展开的表(无技术ID的问题)

该Schema包含2层嵌套数组,展开后应生成3张表:主表、businesspartners表、businesspartner_addresses表。

主表

# 顶层表,依赖some_business_id作为关联键(但无法保证全局唯一)
main_table = df.select("some_business_id", "some_field")

BusinessPartners表

businesspartner_table = (df
                         .withColumn("businesspartners", F.explode(F.col("businesspartners")))
                         .select(
                             "some_business_id",  # 用于关联主表
                             "businesspartners.name", 
                             "businesspartners.profession"
                         ))

BusinessPartner地址表(关联问题凸显)

businesspartner_address_table = (df
                         .withColumn("businesspartners", F.explode(F.col("businesspartners")))
                         .withColumn("addresses", F.explode(F.col("businesspartners.addresses")))
                         .select(
                             "some_business_id",
                             "businesspartners.name",
                             "businesspartners.profession",
                             "addresses.street",
                             "addresses.city",
                         ))

当前用some_business_id关联主表与businesspartners表,但该字段无法保证唯一性;尝试用name+profession组合作为businesspartners表的关联键,也无法确保全局唯一,导致地址表与businesspartners表的关联存在数据歧义,必须引入技术ID解决。曾尝试monotonically_increasing_id()但未找到最优实现方式。


解决方案

关于PySpark工具包

目前没有专门针对PySpark、同时支持自动展开嵌套数据并生成层级关联技术ID的成熟工具包,需要通过自定义递归逻辑实现需求。

实践代码:递归展开嵌套数据并生成关联技术ID

以下代码会自动遍历DataFrame的Schema,递归处理嵌套数组,生成每个层级的独立DataFrame,并添加父/子技术ID用于关联:

import pyspark.sql.functions as F
from pyspark.sql.types import StructType, ArrayType

def generate_tech_id():
    """生成全局唯一技术ID(分布式场景推荐用F.expr("uuid()")替代)"""
    return F.expr("uuid()").alias("tech_id")

def flatten_nested_data(df, parent_id_col=None, parent_table_name="main"):
    """
    递归展开嵌套DataFrame,生成各层级独立表及关联ID
    :param df: 输入的嵌套DataFrame
    :param parent_id_col: 父表的技术ID列名(顶层表为None)
    :param parent_table_name: 父表名称,用于命名子表
    :return: 字典,键为表名,值为对应的DataFrame
    """
    tables = {}
    
    # 1. 处理当前层级的表:提取非数组字段,添加技术ID
    current_fields = [field.name for field in df.schema.fields if not isinstance(field.dataType, ArrayType)]
    current_df = df.select(*current_fields)
    
    if parent_id_col is None:
        # 顶层表:添加自身技术ID
        current_df = current_df.withColumn("tech_id", generate_tech_id())
        tables[parent_table_name] = current_df
    else:
        # 子表:保留父ID,添加自身技术ID
        current_df = current_df.withColumn("parent_tech_id", F.col(parent_id_col)) \
                               .withColumn("tech_id", generate_tech_id())
        tables[parent_table_name] = current_df
    
    # 2. 递归处理数组类型字段
    array_fields = [field for field in df.schema.fields if isinstance(field.dataType, ArrayType)]
    for array_field in array_fields:
        field_name = array_field.name
        # 展开数组
        exploded_df = df.withColumn(field_name, F.explode(F.col(field_name)))
        
        # 提取数组内的结构体字段
        struct_subfields = [f"{field_name}.{subfield.name}" for subfield in array_field.dataType.elementType.fields]
        # 携带父层级ID
        select_cols = [parent_id_col] if parent_id_col else ["tech_id"]
        select_cols.extend(struct_subfields)
        exploded_df = exploded_df.select(*select_cols)
        
        # 子表命名与递归处理
        child_table_name = f"{parent_table_name}_{field_name}"
        child_tables = flatten_nested_data(exploded_df, parent_id_col=select_cols[0], parent_table_name=child_table_name)
        
        # 合并子表结果
        tables.update(child_tables)
    
    return tables

# 使用示例
result_tables = flatten_nested_data(df)

# 查看生成的所有表
for table_name, table_df in result_tables.items():
    print(f"=== 表名: {table_name} ===")
    table_df.show(truncate=False)

代码说明

  1. 技术ID生成:默认用Spark内置的uuid()函数生成全局唯一ID,适合分布式场景,避免Driver端性能瓶颈;也可替换为雪花ID、monotonically_increasing_id()结合row_number()等规则。
  2. 递归逻辑:自动识别Schema中的数组字段,逐层展开并生成子表,每个子表通过parent_tech_id关联父表的tech_id,彻底解决关联键不唯一的问题。
  3. 表命名规则:子表名称基于父表名称+数组字段名,例如主表的businesspartners数组生成main_businesspartners表,其下的addresses数组生成main_businesspartners_addresses表,便于识别层级关系。
  4. Schema兼容性:无需手动编写展开逻辑,自动适配不同市场的Schema差异,满足可扩展需求。

优化建议

  • 若处理超大规模数据,可添加分区逻辑,减少展开操作的计算压力。
  • 可增加字段过滤规则,跳过无需展开的数组或结构体字段。
  • 可将生成的表自动写入Delta Lake,通过技术ID维护表间关联关系。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 16:15:12