PySpark使用explode函数后保留字段名层级的实现方法
解决方案
思路
通过自定义通用函数,自动识别嵌套数组中的结构体字段,在explode后将结构体字段展开,并自动以「父列名.子字段名」的格式命名新列,无需手动指定或重命名。
代码实现
1. 创建测试DataFrame
from pyspark.sql import SparkSession from pyspark.sql.functions import explode, col spark = SparkSession.builder.appName("NestedFieldExpand").getOrCreate() # 模拟输入数据 data = [ { "KTYPE": [{"name": "john", "surname": "elo"}, {"name": "dd", "surname": "ss"}], "MTYPE": [{"name": "dsdsd", "id": "elo"}, {"name": "sdss", "id": "sds22"}] } ] df = spark.createDataFrame(data)
2. 定义通用处理函数
def explode_with_prefix(df, array_col_name): # 对目标数组列执行explode操作 exploded_df = df.withColumn(array_col_name, explode(col(array_col_name))) # 获取该数组列嵌套结构体的所有字段名 struct_fields = exploded_df.select(f"{array_col_name}.*").columns # 逐个展开结构体字段,自动添加父列名作为前缀 for field in struct_fields: new_col_name = f"{array_col_name}.{field}" exploded_df = exploded_df.withColumn(new_col_name, col(f"{array_col_name}.{field}")) # 删除原数组列(此时已变为结构体类型) return exploded_df.drop(array_col_name)
3. 处理目标列并查看结果
# 依次处理KTYPE和MTYPE列 result_df = explode_with_prefix(df, "KTYPE") result_df = explode_with_prefix(result_df, "MTYPE") # 打印结果数据和Schema result_df.show(truncate=False) result_df.printSchema()
效果说明
执行后,DataFrame的列名将自动变为KTYPE.name、KTYPE.surname、MTYPE.name、MTYPE.id,完全符合需求,且无需手动指定每个字段的重命名规则,适配任意嵌套字段数量的数组列。
内容的提问来源于stack exchange,提问作者new_programmer11
相关产品推荐
相关产品推荐

