基于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
相关产品推荐
相关产品推荐

