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

PySpark DataFrame展开JSON数组列失败,请求排查错误原因

问题

我在PySpark DataFrame中有一个名为substitutions的JSON字符串列,该列存储着包含多个对象的数组,希望将数组展开为每行对应一个元素的形式,DataFrame中还有其他列。

DataFrame结构如下:

+--------------------+----------+--------------------+--------------------+----------+--------------------+---------+--------+--------------------+
|           requestid|sourcepage|              cartid|                  tm|        dt|          customerId| usItemId|prefType|       substitutions|
+--------------------+----------+--------------------+--------------------+----------+--------------------+---------+--------+--------------------+
|00-efbedfe05b4482...|  CHECKOUT|808b44cc-1a38-4dd...|2023-04-25 00:07:...|2023-04-25|f1a34e16-a6d0-6f5...|862776084| NO_PREF|{"id":{"productId...|
+--------------------+----------+--------------------+--------------------+----------+--------------------+---------+--------+--------------------+

substitutions列的JSON字符串内容为:

[
    {
        "id": {
            "productId": "2N3UYGUTROQK",
            "usItemId": "32667929"
        },
        "usItemId": "32667929",
        "itemRank": 1,
        "customerChoice": false
    },
    {
        "id": {
            "productId": "2N3UYGUTRHQK",
            "usItemId": "32667429"
        },
        "usItemId": "32667429",
        "itemRank": 2,
        "customerChoice": true
    },
    {
        "id": {
            "productId": "2N3UYGUTRYQK",
            "usItemId": "32667529"
        },
        "usItemId": "32667529",
        "itemRank": 3,
        "customerChoice": false
    },
    {
        "id": {
            "productId": "2N3UYGUTIQK",
            "usItemId": "32667329"
        },
        "usItemId": "32667329",
        "itemRank": 4,
        "customerChoice": false
    },
    {
        "id": {
            "productId": "2N3UYGUTYOQK",
            "usItemId": "32663929"
        },
        "usItemId": "32663929",
        "itemRank": 5,
        "customerChoice": false
    }
]

尝试了以下代码但未得到预期结果:

df.select("*", f.explode(f.from_json("substitutions", MapType(StringType(),StringType()))))

结果:

+--------------------+----------+--------------------+--------------------+----------+--------------------+---------+--------+--------------------+-------+
|           requestid|sourcepage|              cartid|                  tm|        dt|          customerId| usItemId|prefType|       substitutions|entries|
+--------------------+----------+--------------------+--------------------+----------+--------------------+---------+--------+--------------------+-------+
|00-efbedfe05b4482...|  CHECKOUT|808b44cc-1a38-4dd...|2023-04-25 00:07:...|2023-04-25|f1a34e16-a6d0-6f5...|862776084| NO_PREF|[{"id":{"productI...|   null|
+--------------------+----------+--------------------+--------------------+----------+--------------------+---------+--------+--------------------+-------+

请问我哪里出错了?


错误原因

你用错了Schema类型:

  • substitutions是JSON数组,里面每个元素是嵌套结构的对象,不是简单的键值对Map。
  • MapType(StringType(), StringType())只能解析键值都是字符串的JSON对象,无法匹配数组+嵌套对象的结构,导致from_json返回null,后续explode自然也得不到有效数据。

正确解决方案

需要先定义匹配JSON结构的Schema,解析数组后再展开:

步骤1:定义Schema

根据substitutions的JSON结构,用ArrayType包裹StructType,包含所有嵌套字段:

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

# 定义嵌套的id结构
id_schema = StructType([
    StructField("productId", StringType()),
    StructField("usItemId", StringType())
])

# 定义数组元素的结构
substitution_schema = StructType([
    StructField("id", id_schema),
    StructField("usItemId", StringType()),
    StructField("itemRank", IntegerType()),
    StructField("customerChoice", BooleanType())
])

# 最终的数组Schema
full_schema = ArrayType(substitution_schema)

步骤2:解析JSON并展开数组

先解析substitutions列成数组类型,再用explode展开为多行:

from pyspark.sql import functions as f

# 解析JSON列
df_parsed = df.withColumn("substitutions_parsed", f.from_json("substitutions", full_schema))

# 展开数组,同时保留原有的其他列
df_exploded = df_parsed.select("*", f.explode("substitutions_parsed").alias("substitution"))

# 可选:展开嵌套字段,把substitution里的字段拆成单独列
df_final = df_exploded.drop("substitutions", "substitutions_parsed") \
    .select(
        "requestid", "sourcepage", "cartid", "tm", "dt", "customerId", "usItemId", "prefType",
        "substitution.id.productId", "substitution.id.usItemId",
        "substitution.usItemId", "substitution.itemRank", "substitution.customerChoice"
    )

最终效果

展开后每行对应一个substitution元素,原有的其他列会被复制到每行,嵌套的字段也可以单独提取出来。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 03:59:53