如何在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
相关产品推荐
相关产品推荐

