如何基于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") )
动态生成解决方案
实现思路
- 过滤出业务数据中存在的有效字段映射,避免处理无数据的字段
- 将嵌套路径字符串拆分为层级列表,构建树状结构
- 递归遍历树结构,动态生成对应的
F.struct表达式 - 将生成的表达式应用到业务数据表,得到嵌套结构
完整代码实现
# 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
相关产品推荐
相关产品推荐

