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

求通用Python/PySpark脚本:将无规则深度嵌套MongoDB JSON转关系型数据集

递归提取嵌套JSON生成关联关系型数据集(Pandas/PySpark实现)

核心思路

  • 递归遍历JSON字段,识别嵌套对象(dict)和数组(list),拆分为独立关系表
  • 为每个表自动生成唯一主键id,通过parent_id关联父表,保障表间关联关系
  • 针对大数据量:Pandas采用逐记录处理+增量写入,PySpark利用分布式计算规避内存溢出
  • 兼容无固定模式:自动适配不同记录的结构差异,动态生成表和字段

Pandas 实现(适合中小规模数据或分批处理大数据)

代码实现

import pandas as pd
import uuid
from typing import Dict, List

def recursive_flatten(
    data: Dict, 
    parent_id: str = None, 
    parent_table: str = "main", 
    tables: Dict[str, pd.DataFrame] = None
) -> Dict[str, pd.DataFrame]:
    if tables is None:
        tables = {}
    
    # 生成当前记录唯一ID
    current_id = str(uuid.uuid4())
    # 分离普通字段与嵌套字段
    flat_data = {}
    nested_fields = {}
    
    for key, value in data.items():
        if isinstance(value, dict):
            nested_fields[key] = value
        elif isinstance(value, list) and all(isinstance(item, dict) for item in value):
            nested_fields[key] = value
        else:
            flat_data[key] = value
    
    # 添加关联字段
    flat_data["id"] = current_id
    if parent_id:
        flat_data["parent_id"] = parent_id
    
    # 更新当前表
    if parent_table not in tables:
        tables[parent_table] = pd.DataFrame([flat_data])
    else:
        tables[parent_table] = pd.concat([tables[parent_table], pd.DataFrame([flat_data])], ignore_index=True)
    
    # 递归处理嵌套字段
    for nested_key, nested_value in nested_fields.items():
        child_table = f"{parent_table}_{nested_key}"
        if isinstance(nested_value, dict):
            recursive_flatten(nested_value, current_id, child_table, tables)
        elif isinstance(nested_value, list):
            for item in nested_value:
                recursive_flatten(item, current_id, child_table, tables)
    
    return tables

# 示例使用(替换为MongoDB游标遍历即可处理大数据)
sample_data = {
    "name": "John Doe",
    "age": 30,
    "address": {"street": "123 Main St", "city": "Anytown"},
    "orders": [{"order_id": "ORD123", "total": 99.99}, {"order_id": "ORD456", "total": 49.99}]
}

# 生成关联表
tables = recursive_flatten(sample_data)

# 输出所有表(可写入CSV/Parquet到数据湖)
for table_name, df in tables.items():
    print(f"=== 表: {table_name} ===")
    print(df)
    print("\n")

关键说明

  • 用uuid生成全局唯一主键,避免跨记录主键冲突
  • 自动识别嵌套结构,动态生成子表(如main_address、main_orders)
  • 逐记录处理,适配MongoDB游标分批读取场景,避免内存溢出
  • 兼容结构不一致的记录:不同记录的嵌套字段会自动合并到对应表,缺失字段填充NaN

PySpark 实现(适合大规模分布式数据)

代码实现

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, explode, uuid
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DoubleType, ArrayType

def recursive_unroll(df, parent_table="main", parent_id_col="id"):
    tables = {parent_table: df}
    # 识别嵌套字段(结构体/结构体数组)
    nested_cols = [
        (name, dtype) for name, dtype in df.dtypes 
        if dtype.startswith("struct") or dtype.startswith("array<struct")
    ]
    
    for col_name, col_type in nested_cols:
        child_table = f"{parent_table}_{col_name}"
        # 处理数组类型
        if col_type.startswith("array<struct"):
            child_df = df.select(col(parent_id_col).alias("parent_id"), explode(col(col_name)).alias("data"))
            child_df = child_df.select("parent_id", "data.*")
        # 处理结构体类型
        else:
            child_df = df.select(col(parent_id_col).alias("parent_id"), col(col_name).alias("data"))
            child_df = child_df.select("parent_id", "data.*")
        
        # 为子表生成主键
        child_df = child_df.withColumn("id", uuid())
        tables[child_table] = child_df
        
        # 递归处理子表的嵌套字段
        child_tables = recursive_unroll(child_df, child_table, "id")
        tables.update(child_tables)
    
    return tables

# 初始化Spark会话
spark = SparkSession.builder.appName("NestedJsonToRelational").getOrCreate()

# 示例数据(替换为MongoDB读取逻辑:spark.read.format("mongo").option("uri", "mongodb://localhost:27017/your_db.your_collection").load())
sample_schema = StructType([
    StructField("name", StringType()),
    StructField("age", IntegerType()),
    StructField("address", StructType([
        StructField("street", StringType()),
        StructField("city", StringType())
    ])),
    StructField("orders", ArrayType(StructType([
        StructField("order_id", StringType()),
        StructField("total", DoubleType())
    ])))
])

sample_data = [
    {
        "name": "John Doe",
        "age": 30,
        "address": {"street": "123 Main St", "city": "Anytown"},
        "orders": [{"order_id": "ORD123", "total": 99.99}, {"order_id": "ORD456", "total": 49.99}]
    }
]

df = spark.createDataFrame(sample_data, schema=sample_schema)
# 为主表生成主键
df = df.withColumn("id", uuid())

# 生成关联表
tables = recursive_unroll(df)

# 输出所有表(可写入Parquet到数据湖)
for table_name, df in tables.items():
    print(f"=== 表: {table_name} ===")
    df.show(truncate=False)
    print("\n")
    # df.write.mode("overwrite").parquet(f"your-data-lake-path/{table_name}/")

关键说明

  • 依托Spark分布式计算能力,支持TB级数据处理,无内存压力
  • 自动识别嵌套结构,递归展开为关联子表
  • 用uuid()生成分布式唯一主键,保障跨节点唯一性
  • 原生支持写入Parquet/ORC等数据湖存储格式,符合数据湖规范

最佳实践补充

  1. 元数据管理:记录每个表的字段来源、数据类型,便于后续数据治理
  2. 字段冲突处理:若不同层级有同名字段,自动添加前缀(如main_address_street)避免混淆
  3. 增量更新:从MongoDB读取时,基于_id或时间戳增量同步,减少全量扫描开销
  4. 数据验证:添加字段类型校验,处理JSON中非标准数据(如字符串类型的数字)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 22:55:18