如何从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)
代码说明
- 展开嵌套数组:使用
explode函数将JSON中的数组字段展开,每个数组元素对应一行记录。 - 提取目标字段:通过
col函数从展开后的JSON对象中提取url和productNumber,并重命名为统一列名。 - 标记来源:用
lit函数添加source列,区分数据来自哪个原始JSON列。 - 合并数据:使用
unionByName合并三个DataFrame,确保列名匹配,避免因列顺序不同导致的错误。
内容的提问来源于stack exchange,提问作者batman_special
相关产品推荐
相关产品推荐

