求通用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等数据湖存储格式,符合数据湖规范
最佳实践补充
- 元数据管理:记录每个表的字段来源、数据类型,便于后续数据治理
- 字段冲突处理:若不同层级有同名字段,自动添加前缀(如
main_address_street)避免混淆 - 增量更新:从MongoDB读取时,基于
_id或时间戳增量同步,减少全量扫描开销 - 数据验证:添加字段类型校验,处理JSON中非标准数据(如字符串类型的数字)
内容的提问来源于stack exchange,提问作者Miguel Botelho
相关产品推荐
相关产品推荐

