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

如何从PySpark DataFrame嵌套JSON中提取列值并合并结果

PySpark嵌套JSON列提取与合并方案

需求说明

处理名为es_query的PySpark DataFrame,该DataFrame包含r_json、brd_json、vs_json三个嵌套JSON列,需从这三个列中提取url和productNumber字段,将每条提取出的记录作为独立行存入e_result DataFrame,最终合并所有数据到单一DataFrame。

示例数据

// r_json结构及示例
r_json:
{
  "results": [
    {"col1": "Yes", "name": "", "col2": 1, "col3": "76,67 €", "col4": "5,75 €", "productNumber": "B0e28213", "url": "https://www.am"},
    {"col1": "Yes", "name": "", "col2": 1, "col3": "76,67 €", "col4": "5,75 €", "productNumber": "019883", "url": "https://www.am"}
  ]
}

// brd_json结构及示例
brd_json:
{
  "array": [
    {"col1": "Yes", "col2": "https://m.media-a", "col3": null, "col4": "Yes", "col5": "No", "col6": false, "productNumber": "11873628", "rating": "4.1", "url": "https://www.amazon"},
    {"col1": "Yes", "col2": "https://m.media-a", "col3": null, "col4": "Yes", "col5": "No", "col6": false, "productNumber": "001838", "rating": "4.1", "url": "https://www.amazon"}
  ]
}

// vs_json结构及示例
vs_json:
{
  "array": [
    {"col1": "Yes", "col2": "https://m.media-a", "col3": null, "col4": "Yes", "col5": "No", "col6": false, "productNumber": "1212", "rating": "4.1", "url": "https://www.amazon"},
    {"col1": "Yes", "col2": "https://m.media-a", "col3": null, "col4": "Yes", "col5": "No", "col6": false, "productNumber": "2321", "rating": "4.1", "url": "https://www.amazon"}
  ]
}

现有尝试的问题

  • 第一段代码未处理嵌套数组,直接提取url无法获取正确数据;且误将brsd_json重复用于视频源数据提取;最后通过groupBy将URL拼接成字符串,不符合每行一条记录的需求。
  • 第二段代码将每个字段拆分为单独DataFrame,增加了合并复杂度,无需拆分字段,应先展开数组再提取目标字段。

解决方案代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import explode, col, lit

# 初始化SparkSession(如果未初始化)
spark = SparkSession.builder.appName("ExtractJsonData").getOrCreate()

# 处理r_json:展开results数组,提取url和productNumber,标记来源
r_df = es_query.select(
    explode(col("r_json.results")).alias("r_data")
).select(
    col("r_data.url").alias("url"),
    col("r_data.productNumber").alias("productNumber"),
    lit("result").alias("source")
)

# 处理brd_json:展开array数组,提取url和productNumber,标记来源
brd_df = es_query.select(
    explode(col("brd_json.array")).alias("brd_data")
).select(
    col("brd_data.url").alias("url"),
    col("brd_data.productNumber").alias("productNumber"),
    lit("brand").alias("source")
)

# 处理vs_json:展开array数组,提取url和productNumber,标记来源
vs_df = es_query.select(
    explode(col("vs_json.array")).alias("vs_data")
).select(
    col("vs_data.url").alias("url"),
    col("vs_data.productNumber").alias("productNumber"),
    lit("video").alias("source")
)

# 合并三个DataFrame得到最终结果e_result
e_result = r_df.unionByName(brd_df).unionByName(vs_df)

# 查看结果
e_result.show(truncate=False)

代码说明

  1. 展开嵌套数组:使用explode函数将JSON中的数组字段展开,每个数组元素对应一行记录。
  2. 提取目标字段:通过col函数从展开后的JSON对象中提取url和productNumber,并重命名为统一列名。
  3. 标记来源:用lit函数添加source列,区分数据来自哪个原始JSON列。
  4. 合并数据:使用unionByName合并三个DataFrame,确保列名匹配,避免因列顺序不同导致的错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 12:25:18