如何在PySpark中扁平化字符串类型的DataFrame列?
问题分析与解决方案
原代码返回全null的核心原因:你的a列是Python风格的字符串化字典列表(用单引号包裹键值),但Spark的from_json只识别标准JSON格式(双引号),加上用spark.read.json直接读取这种非标准JSON字符串,推断出的Schema完全错误,最终导致解析失败返回null。
正确解决步骤
- 先把
a列的字符串格式转成标准JSON(替换单引号为双引号); - 解析JSON为嵌套结构,再通过两次展开操作处理双层数组,最后提取目标字段。
完整代码示例
from pyspark.sql import functions as F from pyspark.sql.types import ArrayType, StructType, StructField, StringType, LongType # 1. 修复a列的JSON格式:把单引号替换成双引号 df_clean = df.withColumn("a_json", F.regexp_replace(F.col("a"), "'", '"')) # 2. 手动定义匹配数据结构的Schema(比自动推断更稳定) item_schema = StructType([ StructField("npi", ArrayType(LongType())), StructField("tin", StructType([ StructField("type", StringType()), StructField("value", StringType()) ])) ]) array_schema = ArrayType(item_schema) # 3. 解析标准JSON列 df_parsed = df_clean.withColumn("a_parsed", F.from_json(F.col("a_json"), array_schema)) # 4. 展开外层的字典数组(每个字典拆成一行) df_explode_outer = df_parsed.select("b", F.explode("a_parsed").alias("a_item")) # 5. 展开内层的npi数组,同时提取tin的字段 df_final = df_explode_outer.select( F.explode(F.col("a_item.npi")).alias("a_npi"), F.col("a_item.tin.type").alias("a_tin_type"), F.col("a_item.tin.value").alias("a_tin_value"), F.col("b") ) # 查看最终结果 df_final.show()
关键说明
- 替换单引号是核心前提,否则Spark无法识别非标准JSON;
- 手动定义Schema可以避免自动推断时因数据异常导致的格式错误;
- 两次
explode分别处理双层数组,实现完全扁平化的目标输出。
内容的提问来源于stack exchange,提问作者Shri Tech1404
相关产品推荐
相关产品推荐

