PySpark将JSON字符串列拆解为多行多列的实现方案问询
Spark 解析JSON字符串列并展开数组实现方案
场景1:已调整JSON仅保留response数组(推荐)
前置要求:先确保返回的JSON为标准格式,修正多余符号、给所有键加双引号,最终格式示例:[{"to":"Sam", "position":"guard"}, {"to":"John", "position":"center"}, {"to":"Andrew", "position":"forward"}]
完整实现代码如下:
from pyspark.sql import functions as F from pyspark.sql.types import ArrayType, StructType, StructField, StringType # 定义response数组对应的schema response_schema = ArrayType( StructType([ StructField("to", StringType(), nullable=True), StructField("position", StringType(), nullable=True) ]) ) result_df = df \ # 解析JSON字符串为结构化数组 .withColumn("response_arr", F.from_json("col1", response_schema)) \ # 展开数组生成多行数据 .withColumn("item", F.explode("response_arr")) \ # 提取目标字段并重命名 .select( "col2", F.col("item.to").alias("col3"), F.col("item.position").alias("col4") )
执行result_df.show()即可得到预期输出。
场景2:使用原始完整JSON(包含original和response字段)
如果不调整返回的JSON格式,直接解析原始内容的话,schema需要匹配完整JSON结构,实现代码如下:
from pyspark.sql import functions as F from pyspark.sql.types import ArrayType, StructType, StructField, StringType, DoubleType # 定义完整JSON对应的schema full_schema = ArrayType( StructType([ StructField("original", StructType([ StructField("ranking", DoubleType(), nullable=True), StructField("input", StringType(), nullable=True) ]), nullable=True), StructField("response", ArrayType( StructType([ StructField("to", StringType(), nullable=True), StructField("position", StringType(), nullable=True) ]) ), nullable=True) ]) ) result_df = df \ .withColumn("parsed_data", F.from_json("col1", full_schema)) \ # 取数组首个元素的response字段 .withColumn("response_arr", F.col("parsed_data")[0]["response"]) \ .withColumn("item", F.explode("response_arr")) \ .select( "col2", F.col("item.to").alias("col3"), F.col("item.position").alias("col4") )
注意事项
- 必须保证解析的JSON为标准格式,所有键用双引号包裹,不存在多余的括号、逗号等语法错误。如果原始JSON是无引号的类字典格式,可以先用
regexp_replace函数做预处理,例如把=替换为:、给键名补双引号。 - 可根据实际字段类型调整schema中的数据类型,比如ranking为数值类时可对应设置为DoubleType/IntegerType。
内容的提问来源于stack exchange,提问作者dcrowley01
相关产品推荐
相关产品推荐

