如何在PySpark中提取字符串型JSON数据并拆分为多行?
PySpark解析JSON字符串列并生成多行记录
你可以通过以下步骤实现需求:
- 定义JSON结构Schema:明确嵌套JSON的结构,为
from_json函数提供解析依据 - 解析JSON字符串:将字符串类型的
response列转换为结构化数据 - 展开数组:用
explode函数把数组中的每个元素拆分为单独行 - 提取字段并转换类型:从展开后的结构中提取目标字段,同时完成类型转换
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql.functions import from_json, explode, col from pyspark.sql.types import StructType, StructField, StringType, ArrayType, IntegerType # 初始化SparkSession spark = SparkSession.builder.appName("JSONParse").getOrCreate() # 创建示例输入DataFrame data = [("""{"Customer Data": [{"id": "5", "name": "Jony"}, {"id": "10", "name": "Jems"}]}""",)] df = spark.createDataFrame(data, ["response"]) # 定义JSON的Schema schema = StructType([ StructField("Customer Data", ArrayType( StructType([ StructField("id", StringType()), StructField("name", StringType()) ]) )) ]) # 解析JSON、展开数组并提取字段 result_df = df.withColumn("parsed_json", from_json(col("response"), schema)) \ .select(explode(col("parsed_json.Customer Data")).alias("customer")) \ .select( col("customer.id").cast(IntegerType()).alias("id"), col("customer.name").alias("name") ) # 查看结果 result_df.show()
代码说明
- Schema定义:根据你的JSON结构,定义包含数组的StructType,确保
from_json能精准解析字符串列。 - JSON解析:
from_json(col("response"), schema)将字符串列转换为结构化的parsed_json列。 - 数组展开:
explode函数把Customer Data数组中的每个客户对象拆分为独立行,实现一行转多行的效果。 - 字段处理:通过
cast(IntegerType())将字符串类型的id转为整数,同时提取name字段并设置别名。
执行后会输出目标结果:
+---+----+ | id|name| +---+----+ | 5|Jony| | 10|Jems| +---+----+
内容的提问来源于stack exchange,提问作者Jaykumar Patel
相关产品推荐
相关产品推荐

