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

如何在PySpark DataFrame中添加嵌套JSON的notebook_path字段?

问题描述

现有如下Spark代码,用于解析JSON字符串并生成包含job_id和taskname的DataFrame:

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

spark = SparkSession.builder.getOrCreate()

s = '{"job_id":"123","settings":{"task":[{"taskname":"task1","notebook_task":{"notebook_path":"path1"}},{"taskname":"task2","notebook_task":{"notebook_path":"path2"}}]}}'

schema = schema_of_json(lit(s))

result_df = (
    spark.createDataFrame([s], "string")
    .select(from_json(col("value"), schema).alias("data"))
    .select("data.job_id", explode("data.settings.task.taskname").alias("taskname"))
    )

result_df.show()

运行后输出结果:

+------+--------+
|job_id|taskname|
+------+--------+
|   123|   task1|
|   123|   task2|
+------+--------+

需求:需要将notebook_path字段添加到该DataFrame中,希望避免创建额外DataFrame再通过job_id关联的方案,寻求更优解。

最优解决方案

无需拆分后关联,直接先 explode 整个task数组对象,再从每个task元素中提取所需字段即可。这种方式一次explode就能同时获取多个目标字段,逻辑更简洁高效。

修改后的完整代码:

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

spark = SparkSession.builder.getOrCreate()

s = '{"job_id":"123","settings":{"task":[{"taskname":"task1","notebook_task":{"notebook_path":"path1"}},{"taskname":"task2","notebook_task":{"notebook_path":"path2"}}]}}'

schema = schema_of_json(lit(s))

result_df = (
    spark.createDataFrame([s], "string")
    .select(from_json(col("value"), schema).alias("data"))
    # 对整个task数组执行explode,得到单个task对象行
    .select("data.job_id", explode("data.settings.task").alias("task"))
    # 从task对象中提取需要的字段
    .select(
        "job_id",
        col("task.taskname").alias("taskname"),
        col("task.notebook_task.notebook_path").alias("notebook_path")
    )
)

result_df.show()

运行后输出结果:

+------+--------+-------------+
|job_id|taskname|notebook_path|
+------+--------+-------------+
|   123|   task1|        path1|
|   123|   task2|        path2|
+------+--------+-------------+

核心逻辑说明

  • 原代码中explode("data.settings.task.taskname")仅对task数组内的taskname字段做拆分,只能获取单个字段;
  • 改为explode("data.settings.task")后,会将数组中的每个完整task对象拆分为独立行,后续可直接从task对象中提取任意嵌套字段(如taskname、notebook_task.notebook_path),完全不需要关联操作,性能和代码可读性更优。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 02:57:21