如何在PySpark DataFrame中将嵌套字典列表字符串拆分为独立列
处理PySpark中嵌套字典列表字符串列的实现方法
步骤1:将非标准字符串转换为合法JSON格式
原列的字符串格式不符合JSON规范,需要先做批量替换处理:
- 给所有键名(如
ser、cos)添加双引号 - 把
=替换为JSON的键值分隔符: - 保留无引号的
null为JSON格式的null
使用regexp_replace实现:
from pyspark.sql import functions as F # 假设原DataFrame名为df,目标列名为data_col df = df.withColumn("json_str", F.regexp_replace("data_col", r'(\w+)=', r'"\1":')) \ .withColumn("json_str", F.regexp_replace("json_str", r'([{,]\s*)(\w+)', r'\1"\2"'))
步骤2:解析JSON数组为Struct类型
先定义匹配数据结构的Schema,再用from_json解析字符串:
from pyspark.sql.types import ArrayType, StructType, StructField, StringType # 定义嵌套结构的Schema element_schema = StructType([ StructField("ser", StructType([ StructField("cos", StringType()), StructField("mgse", StringType()), StructField("bd", StringType()) ])), StructField("ap", StructType([ StructField("ncd", StringType()), StructField("ccd", StringType()), StructField("scd2", StringType()), StructField("cos", StringType()), StructField("pgse", StringType()), StructField("pcd", StringType()), StructField("nar", StringType()) ])), StructField("ep", StructType([ StructField("eptd", StringType()), StructField("ept", StringType()) ])) ]) df = df.withColumn("parsed_data", F.from_json(F.col("json_str"), ArrayType(element_schema)))
步骤3:展开数组为独立行
用explode把数组中的每个元素拆分成单独的行:
df = df.withColumn("exploded_data", F.explode("parsed_data"))
步骤4:拆分嵌套字段为独立列
依次提取ser、ap、ep下的子字段:
# 拆分ser下的字段 df = df.withColumn("ser_cos", F.col("exploded_data.ser.cos")) \ .withColumn("ser_mgse", F.col("exploded_data.ser.mgse")) \ .withColumn("ser_bd", F.col("exploded_data.ser.bd")) # 拆分ap下的字段 df = df.withColumn("ap_ncd", F.col("exploded_data.ap.ncd")) \ .withColumn("ap_ccd", F.col("exploded_data.ap.ccd")) \ .withColumn("ap_scd2", F.col("exploded_data.ap.scd2")) \ .withColumn("ap_cos", F.col("exploded_data.ap.cos")) \ .withColumn("ap_pgse", F.col("exploded_data.ap.pgse")) \ .withColumn("ap_pcd", F.col("exploded_data.ap.pcd")) \ .withColumn("ap_nar", F.col("exploded_data.ap.nar")) # 拆分ep下的字段 df = df.withColumn("ep_eptd", F.col("exploded_data.ep.eptd")) \ .withColumn("ep_ept", F.col("exploded_data.ep.ept")) # 清理中间列 df = df.drop("data_col", "json_str", "parsed_data", "exploded_data")
补充说明
- 如果字段存在不确定性,可尝试用
schema_of_json自动推断Schema,但稳定性不如手动定义,建议优先明确Schema - 原数据中的空字典(如示例中的
ap={})拆分后对应字段会显示为null,符合缺失值处理逻辑
内容的提问来源于stack exchange,提问作者Bab
相关产品推荐
相关产品推荐

