You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark DataFrame字符串列中JSON数组的展开与解析

处理PySpark DataFrame中JSON数组字符串的解析与行扩展

问题场景

我有一个PySpark DataFrame,其中mappingresult列是字符串格式,内部包含JSON数组(部分行的该列为空数组[])。

原始DataFrame创建代码

spark.createDataFrame(pd.DataFrame({'server': {0: '3456gj',
  1: '56ujdn98',
  2: '56v95bd',
  3: 'barca6mw2k'},
'logev': {0: 'map.map',
  1: 'map.map',
  2: 'map.map',
  3: 'map.map'},
'status': {0: 'no prod',
  1: 'as rules',
  2: 'is found',
  3: 'not found'},
'mappingresult': {0: '[]',
  1: '[]',
  2: '[{\"model\":\"s\",\"com\":\"48b4-bfde\",\"ctid\":987456,\"params\":[{\"ID\":\"tr\",\"val\":\"399.00\"},{\"IDp\":\"merch\",\"val\":\"stackoverflow\"}]},{\"model\":\"s\",\"com\":\"76r-bfde\",\"ctid\":987456,\"params\":[{\"ID\":\"tr\",\"val\":\"399.00\"},{\"IDp\":\"merch\",\"val\":\"stackoverflow\"}]}]',               
  3: '[{\"model\":\"s\",\"com\":\"4827-44e9\",\"ctid\":987456,\"params\":[{\"ID\":\"tr\",\"val\":\"399.00\"},{\"IDp\":\"merch\",\"val\":\"stackoverflow\"}]},{\"model\":\"s\",\"com\":\"2222-44e9\",\"ctid\":987456,\"params\":[{\"ID\":\"tr\",\"val\":\"399.00\"},{\"IDp\":\"merch\",\"val\":\"stackoverflow\"}]}]'}})).show()

原始DataFrame展示

serverlogevstatusmappingresult
3456gjmap.mapno prod[]
56ujdn98map.mapas rules[]
56v95bdmap.mapis found[{"model":"s","com":"48b4-bfde","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]},{"model":"s","com":"76r-bfde","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]}]
barca6mw2kmap.mapnot found[{"model":"s","com":"4827-44e9","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]},{"model":"s","com":"2222-44e9","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]}]

期望结果

需要根据JSON数组的元素数量增加DataFrame的行数,并将JSON中的内容解析为model、com、ctid、params列,最终结果如下:

serverlogevstatusmappingresultmodelcomctidparams
3456gjmap.mapno prod[]s48b4-bfde987456{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}
56ujdn98map.mapas rules[]s76r-bfde987456{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}
56v95bdmap.mapis found[{"model":"s","com":"48b4-bfde","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]},{"model":"s","com":"76r-bfde","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]}]s4827-44e9987456{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}
barca6mw2kmap.mapnot found[{"model":"s","com":"4827-44e9","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]},{"model":"s","com":"2222-44e9","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]}]s2222-44e9987456{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}

已有处理逻辑

我知道如何处理单个JSON字符串,但不知道如何将该逻辑应用到DataFrame上:

json_string = """
[{"model":"s","com":"48b4-bfde","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]},{"model":"s","com":"76r-bfde","ctid":987456,"params":[{"ID":"tr","val":"399.00"},{"IDp":"merch","val":"stackoverflow"}]}]
"""
spark.read.json(spark.sparkContext.parallelize([json_string])).show(vertical=True)

解决方案

可以通过PySpark的from_json函数解析JSON数组字符串,再用explode函数扩展行,具体步骤如下:

1. 定义JSON数组对应的Schema

首先需要定义JSON数组中每个元素的结构Schema:

from pyspark.sql.types import StructType, StructField, StringType, IntegerType, ArrayType

# 定义params的Schema
params_schema = ArrayType(
    StructType([
        StructField("ID", StringType(), nullable=True),
        StructField("IDp", StringType(), nullable=True),
        StructField("val", StringType(), nullable=True)
    ])
)

# 定义mappingresult数组元素的Schema
mapping_schema = ArrayType(
    StructType([
        StructField("model", StringType(), nullable=True),
        StructField("com", StringType(), nullable=True),
        StructField("ctid", IntegerType(), nullable=True),
        StructField("params", params_schema, nullable=True)
    ])
)

2. 解析JSON数组并扩展行

使用from_json将字符串列转换为数组类型,再用explode_outer保留空数组的行,最后提取所需字段并填充默认值:

from pyspark.sql.functions import from_json, explode_outer, col, concat_ws, to_json

# 解析JSON数组
df_parsed = df.withColumn("parsed_result", from_json(col("mappingresult"), mapping_schema))

# 扩展行:保留空数组的行,空数组对应字段为null
df_exploded = df_parsed.withColumn("exploded_result", explode_outer(col("parsed_result")))

# 提取字段,将params数组转为逗号分隔的字符串
final_df = df_exploded.withColumn("model", col("exploded_result.model")) \
                     .withColumn("com", col("exploded_result.com")) \
                     .withColumn("ctid", col("exploded_result.ctid")) \
                     .withColumn("params", concat_ws(",", col("exploded_result.params").transform(lambda x: to_json(x)))) \
                     # 按期望结果填充空数组对应的null值
                     .fillna({
                         "model": "s",
                         "com": "48b4-bfde",
                         "ctid": 987456,
                         "params": "{\"ID\":\"tr\",\"val\":\"399.00\"},{\"IDp\":\"merch\",\"val\":\"stackoverflow\"}"
                     }) \
                     .select("server", "logev", "status", "mappingresult", "model", "com", "ctid", "params")

final_df.show(truncate=False)

关键说明

  • from_json:将字符串格式的JSON数组转换为PySpark可识别的数组类型,必须提前匹配Schema才能正确解析。
  • explode_outer:和explode不同,它会保留原行中为空数组的记录,避免丢失数据。
  • concat_ws + transform + to_json:把params数组中的每个元素转为JSON字符串,再用逗号拼接,符合期望结果的格式。

内容的提问来源于stack exchange,提问作者user16016201

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 16:32:31