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

如何基于PySpark动态构建指定结构的嵌套JSON

PySpark 基于映射表动态生成嵌套JSON结构

问题背景

现有两个Spark DataFrame:

  • 映射表:存储源字段名(data列)与嵌套JSON层级路径(action列)的对应关系
  • 业务数据表:包含实际业务数据

需要基于映射表动态生成符合指定结构的嵌套JSON,替代手动编写F.struct的静态实现方式。

现有数据示例

映射表构建代码

from pyspark.sql import SparkSession
from pyspark.sql import functions as F

spark = SparkSession.builder.appName('sparkdf').getOrCreate()

a1=["DA_STinf","DA_Stinf_NA","DA_Stinf_city","DA_Stinf_NA_ID","DA_Stinf_NA_ID_GRANT","DA_country"]
a2=["data.studentinfo","data.studentinfo.name","data.studentinfo.city","data.studentinfo.name.id","data.studentinfo.name.id.grant","data.country"]
columns = ["data","action"]

df = spark.createDataFrame(zip(a1, a2), columns)

业务数据表构建代码

# 业务数据
a1=["Pune"]
a2=["YES"]
a3=["India"]
col=["DA_Stinf_city","DA_Stinf_NA_ID_GRANT","DA_country"]
data=spark.createDataFrame(zip(a1, a2,a3), col)

预期嵌套JSON结果

{
    "data": {
        "studentinfo": {
            "city": "Pune",
            "name": {
                "id": {
                    "grant": "YES"
                }
            }
        },
        "country": "India"
    }
}

静态手动实现(参考)

以下是手动编写F.struct的实现方式,但仅适用于固定结构,无法动态适配映射表变化:

data.select(        
    F.struct(
        F.struct(
                F.col("DA_Stinf_city").alias("city"),
                F.struct(
                    F.struct(F.col("DA_Stinf_NA_ID_GRANT")).alias("id")
                    ).alias("name"),
        ).alias("studentinfo"),
        F.col("DA_country").alias("country")
    ).alias("data")
)

动态生成解决方案

实现思路

  1. 过滤出业务数据中存在的有效字段映射,避免处理无数据的字段
  2. 将嵌套路径字符串拆分为层级列表,构建树状结构
  3. 递归遍历树结构,动态生成对应的F.struct表达式
  4. 将生成的表达式应用到业务数据表,得到嵌套结构

完整代码实现

# 1. 过滤出业务数据中存在的有效映射
valid_mappings = df.filter(F.col("data").isin(data.columns)).collect()

# 2. 构建层级树结构
tree = {}
for row in valid_mappings:
    src_col = row["data"]
    # 拆分路径为层级列表
    path_parts = row["action"].split(".")
    current_node = tree
    # 遍历路径层级,创建嵌套节点
    for part in path_parts[:-1]:
        if part not in current_node:
            current_node[part] = {}
        current_node = current_node[part]
    # 叶子节点关联源字段
    current_node[path_parts[-1]] = src_col

# 3. 递归生成struct表达式
def build_struct(node):
    struct_fields = []
    for key, value in node.items():
        if isinstance(value, dict):
            # 嵌套节点,递归生成子struct
            sub_struct = build_struct(value)
            struct_fields.append(sub_struct.alias(key))
        else:
            # 叶子节点,直接关联源字段
            struct_fields.append(F.col(value).alias(key))
    return F.struct(*struct_fields)

# 4. 应用到业务数据生成嵌套结构
result = data.select(build_struct(tree).alias("data"))

# 查看结果
print(result.toJSON().collect()[0])

代码说明

  • 树结构构建:将每个路径拆解为层级节点,把源字段挂载到对应的叶子节点,确保嵌套关系正确
  • 递归生成struct:自动识别嵌套层级,对每个节点生成对应的struct,无需手动编写多层嵌套代码
  • 灵活性:无论映射表的嵌套层级多深,只要路径格式正确,都能自动适配生成对应结构

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 17:40:24