如何用PySpark递归展平多嵌套类型复杂JSON并处理空值
递归展平包含嵌套Struct、Array、Map的PySpark DataFrame(支持空值处理)
当然可以实现递归处理所有嵌套类型的PySpark DataFrame展平函数,下面是完整实现,支持Struct、Array、Map的多层嵌套,同时处理空Struct、空Array、空Map等空值场景:
完整代码实现
from pyspark.sql import functions as F from pyspark.sql.types import StructType, ArrayType, MapType def flatten_complex_df(df, parent_col=""): selected_cols = [] for col_name, col_type in df.dtypes: # 拼接完整列名,避免嵌套字段名冲突 full_col_name = f"{parent_col}_{col_name}" if parent_col else col_name current_col = F.col(col_name) if not parent_col else F.col(f"{parent_col}.{col_name}") # 处理Struct类型:递归展平内部字段 if col_type.startswith("struct"): # 空Struct转为默认空结构,防止递归报错 struct_schema = df.schema[col_name].dataType empty_struct = F.struct(*[F.lit(None).alias(field.name) for field in struct_schema.fields]) safe_col = F.coalesce(current_col, empty_struct) nested_df = safe_col.alias(col_name).select(col_name + ".*") nested_cols = flatten_complex_df(nested_df, parent_col=full_col_name) selected_cols.extend(nested_cols) # 处理Array类型:先展开数组,再递归展平内部元素 elif col_type.startswith("array"): # 用explode_outer保留空数组对应的行 exploded_col = F.explode_outer(current_col).alias(f"{col_name}_elem") elem_type = df.schema[col_name].dataType.elementType if isinstance(elem_type, (StructType, ArrayType, MapType)): # 数组元素是复杂类型,递归展平 nested_df = exploded_col.alias(f"{col_name}_elem").select(f"{col_name}_elem.*") nested_cols = flatten_complex_df(nested_df, parent_col=f"{full_col_name}_elem") # 添加数组索引,区分同数组的不同元素 selected_cols.append(F.monotonically_increasing_id().alias(f"{full_col_name}_idx")) selected_cols.extend(nested_cols) else: # 数组元素是简单类型,直接重命名列 selected_cols.append(exploded_col.alias(full_col_name)) # 处理Map类型:转为键值对结构后展平 elif col_type.startswith("map"): # 空Map转为空结构,避免后续操作报错 safe_col = F.coalesce(current_col, F.create_map()) value_type = df.schema[col_name].dataType.valueType if isinstance(value_type, (StructType, ArrayType, MapType)): # Map的Value是复杂类型,先拆分为键值对Struct再展平 map_entries = F.map_entries(safe_col).alias(f"{col_name}_entry") exploded_entry = F.explode_outer(map_entries).alias(f"{col_name}_entry") key_col = exploded_entry["key"].alias(f"{full_col_name}_key") value_col = exploded_entry["value"].alias(f"{col_name}_value") nested_df = value_col.select(f"{col_name}_value.*") nested_cols = flatten_complex_df(nested_df, parent_col=f"{full_col_name}_value") selected_cols.append(key_col) selected_cols.extend(nested_cols) else: # Map的Value是简单类型,直接拆分为键和值两列 map_entries = F.map_entries(safe_col).alias(f"{col_name}_entry") exploded_entry = F.explode_outer(map_entries).alias(f"{col_name}_entry") selected_cols.append(exploded_entry["key"].alias(f"{full_col_name}_key")) selected_cols.append(exploded_entry["value"].alias(f"{full_col_name}_value")) # 简单类型:直接重命名列 else: selected_cols.append(current_col.alias(full_col_name)) return selected_cols # 使用示例 # flattened_df = original_json_df.select(flatten_complex_df(original_json_df))
核心特性说明
- 全类型递归处理:自动识别并递归展平Struct、Array、Map类型,直到所有字段都是简单类型(如string、int、double等)
- 空值安全处理:
- 空Struct:通过
coalesce转为默认空Struct,避免递归时因null值报错 - 空Array:使用
explode_outer替代explode,保留原数据行而不丢失 - 空Map:转为空Map结构,确保
map_entries等操作正常执行
- 空Struct:通过
- 列名冲突防护:嵌套字段名用下划线
_拼接,数组元素加_elem后缀,Map键值分别加_key/_value后缀,避免列名重复 - 动态适配结构:无需手动指定嵌套层级,不管JSON结构多复杂,都能自动适配处理
内容的提问来源于stack exchange,提问作者harshith
相关产品推荐
相关产品推荐

