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

基于DataFrame动态构建JSON Schema与Spark数据提取结构

基于DataFrame动态生成JSON Schema与Spark提取语句

首先是用于生成逻辑的DataFrame构建代码:

import pandas as pd
import numpy as np

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"]
a3 = [np.NaN, np.NaN, "StringType", np.NaN, "BoolType", "StringType"]

d1 = pd.DataFrame(list(zip(a1, a2, a3)), columns=['data', 'action', 'datatype'])

1. 生成符合StructType([StructField(Column_name,Datatype,True)])格式的JSON Schema

通过拆解action列的层级路径,构建树形结构后递归生成嵌套Schema,实现代码如下:

from pyspark.sql.types import StructType, StructField, StringType, BooleanType

# 类型映射字典
type_mapping = {
    "StringType": StringType(),
    "BoolType": BooleanType()
}

# 构建嵌套结构树
schema_tree = {}
for _, row in d1.iterrows():
    path_parts = row['action'].split('.')
    current_node = schema_tree
    # 遍历路径除最后一个节点的部分
    for part in path_parts[:-1]:
        if part not in current_node:
            current_node[part] = {"children": {}}
        current_node = current_node[part]["children"]
    # 处理最后一个字段的类型
    field_name = path_parts[-1]
    current_node[field_name] = type_mapping.get(row['datatype'], StringType()) if pd.notna(row['datatype']) else None

# 递归生成StructField
def build_struct(node):
    fields = []
    for name, value in node.items():
        if isinstance(value, dict):
            # 生成嵌套StructType
            nested_struct = StructType(build_struct(value["children"]))
            fields.append(StructField(name, nested_struct, True))
        else:
            # 生成基础类型字段,未指定类型默认用StringType
            field_type = value if value is not None else StringType()
            fields.append(StructField(name, field_type, True))
    return fields

# 生成根Schema
root_schema = StructType(build_struct(schema_tree))

最终输出的JSON Schema

{
  "type": "struct",
  "fields": [
    {
      "name": "data",
      "type": {
        "type": "struct",
        "fields": [
          {
            "name": "studentinfo",
            "type": {
              "type": "struct",
              "fields": [
                {
                  "name": "name",
                  "type": {
                    "type": "struct",
                    "fields": [
                      {
                        "name": "id",
                        "type": {
                          "type": "struct",
                          "fields": [
                            {
                              "name": "grant",
                              "type": "boolean",
                              "nullable": true
                            }
                          ]
                        },
                        "nullable": true
                      }
                    ]
                  },
                  "nullable": true
                },
                {
                  "name": "city",
                  "type": "string",
                  "nullable": true
                }
              ]
            },
            "nullable": true
          },
          {
            "name": "country",
            "type": "string",
            "nullable": true
          }
        ]
      },
      "nullable": true
    }
  ]
}

2. 生成符合F.struct(F.col(column_name)).alias(json_expected_name)格式的Spark数据提取语句

同样基于路径拆解构建表达式树,递归生成嵌套的struct提取逻辑,实现代码如下:

from pyspark.sql import functions as F

# 构建提取表达式树
expr_tree = {}
for _, row in d1.iterrows():
    path_parts = row['action'].split('.')
    alias_name = row['data']
    current_node = expr_tree
    # 遍历路径到倒数第二节点
    for part in path_parts[:-1]:
        if part not in current_node:
            current_node[part] = {"children": {}, "alias": None}
        current_node = current_node[part]["children"]
    # 绑定最后一个字段与别名
    field_name = path_parts[-1]
    current_node[field_name] = {"alias": alias_name}

# 递归生成struct表达式
def build_expr(node):
    struct_fields = []
    for name, value in node.items():
        if isinstance(value, dict):
            if value.get("alias"):
                # 叶子节点生成带别名的列表达式
                struct_fields.append(F.col(name).alias(value["alias"]))
            else:
                # 嵌套节点递归生成struct
                nested_expr = build_expr(value["children"])
                struct_fields.append(F.struct(*nested_expr).alias(name))
    return struct_fields

# 生成最终提取语句
extract_expr = F.struct(*build_expr(expr_tree))

最终输出的Spark提取语句(格式化后)

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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 16:01:28