Spark Python自动扁平化JSON:仅展开值为String类型的字段
Spark Python 自动扁平化JSON(仅展开String类型字段)
需求说明
无需手动指定字段,自动遍历JSON嵌套结构,仅将值为String类型的字段展开到顶层,非String类型的嵌套/字段保留原有结构。
实现方案
通过递归遍历DataFrame Schema识别所有String类型字段的完整路径,再基于这些路径生成选择表达式,实现自动扁平化。
步骤1:递归收集String类型字段路径
编写函数遍历Schema,递归记录所有String类型字段的全路径(含嵌套层级):
from pyspark.sql.types import StructType, StringType from pyspark.sql.functions import col def get_string_field_paths(schema, parent_path=""): """递归遍历Schema,返回所有String类型字段的完整路径""" paths = [] for field in schema.fields: current_path = f"{parent_path}.{field.name}" if parent_path else field.name if isinstance(field.dataType, StructType): # 递归处理嵌套结构体 paths.extend(get_string_field_paths(field.dataType, current_path)) elif isinstance(field.dataType, StringType): paths.append(current_path) return paths
步骤2:生成扁平化DataFrame
基于收集到的路径,构建选择表达式,将String字段展开到顶层(可自定义命名规则,比如用下划线替换点分隔符):
# 假设df是你的原始JSON DataFrame string_field_paths = get_string_field_paths(df.schema) # 构建选择表达式:保留非结构体的顶层字段 + 展开的String字段 select_expr = [] # 保留顶层非结构体字段(非String/非Struct的字段直接保留) for field in df.schema.fields: if not isinstance(field.dataType, StructType): select_expr.append(col(field.name)) # 添加展开的String字段,并重命名避免点号问题 for path in string_field_paths: # 将路径中的点替换为下划线作为新字段名 alias_name = path.replace(".", "_") select_expr.append(col(path).alias(alias_name)) # 生成最终扁平化DataFrame flattened_df = df.select(*select_expr)
示例验证
原始数据样例
{ "id": 101, "details": { "username": "greencolor", "role": "developer", "stats": { "join_date": "2023-01-01", "contributions": 150 } }, "is_verified": true }
预期输出结构
| id | is_verified | details_username | details_role | details_stats_join_date |
|---|---|---|---|---|
| 101 | true | greencolor | developer | 2023-01-01 |
(注:原details.stats.contributions为int类型,保留在原嵌套结构中;若需移除原嵌套结构,可修改选择表达式排除Struct类型字段)
可选调整
- 移除原嵌套结构:如果不需要保留原始的Struct字段,在构建
select_expr时直接排除所有Struct类型的顶层字段即可。 - 自定义字段命名:可修改
alias_name的生成规则,比如保留点号或使用其他分隔符。
内容的提问来源于stack exchange,提问作者Greencolor
相关产品推荐
相关产品推荐

