如何将类JSON记录转换为Dataset并映射为自定义类对象?
问题
我有一个存储单行JSON记录的文件,示例记录如下(实际每条记录为单行):
{"name": "pim pom", "types": "amy \n klim\nshining rock(ABC)\nflying\nchanning", "url": "http://doingrock.com", "image": "http://static.doingrock.com/rockisland.jpg", "pullTime": "PT3AM", "rockHeight": "8", "dateLive": "2010-10-14", "hitTime": "PT8PM", "desc": "Amazing view"}
我希望读取所有记录并转换为Dataset,将JSON记录存储为类对象,以便通过类属性访问数据进行后续计算。
当前实现代码如下:
schema = StructType([StructField("name", StringType(), True), StructField("ingredients", StringType(), True), StructField("url", StringType(), True), StructField("images", StringType(), True), StructField("pullTime", StringType(), True), StructField("rockHeight", StringType(), True), StructField("dateLive", StringType(), True), StructField("hitTime", StringType(), True), StructField("desc", StringType(), True)]) dataset = spark.sparkContext.textFile('D:/pythonwork/FUN/input/*') dataDicts = dataset.toDF().select(from_json(dataset, schema).alias("dicts"))
目前代码生成的是字符串类型JSON的RDD,请问能否创建类结构将JSON字符串存储为类对象,进而生成Dataset
解决方案
当然可以创建类结构来存储JSON数据,生成强类型的Dataset,以下是具体实现步骤:
1. 修正Schema字段匹配问题
你的现有Schema存在字段名不匹配的问题:JSON中的types对应代码里的ingredients,image对应images,这会导致JSON解析失败,先修正Schema:
from pyspark.sql.types import StructType, StructField, StringType schema = StructType([ StructField("name", StringType(), True), StructField("types", StringType(), True), # 修正为JSON中的字段名 StructField("url", StringType(), True), StructField("image", StringType(), True), # 修正为JSON中的字段名 StructField("pullTime", StringType(), True), StructField("rockHeight", StringType(), True), StructField("dateLive", StringType(), True), StructField("hitTime", StringType(), True), StructField("desc", StringType(), True) ])
2. 定义对应的数据类
在PySpark中,使用dataclasses定义与JSON结构匹配的数据类,Spark可自动生成对应的Encoder来支持强类型Dataset操作:
from dataclasses import dataclass from pyspark.sql import Encoders @dataclass class RockRecord: name: str types: str url: str image: str pullTime: str rockHeight: str dateLive: str hitTime: str desc: str # 生成数据类对应的Encoder rock_encoder = Encoders.product(RockRecord)
3. 读取并转换为强类型Dataset
读取文件后,先解析JSON为结构化DataFrame,再转换为对应数据类的Dataset:
from pyspark.sql.functions import from_json # 读取文本文件转为DataFrame text_df = spark.sparkContext.textFile('D:/pythonwork/FUN/input/*').toDF("json_str") # 解析JSON字符串为结构化数据 parsed_df = text_df.select(from_json(text_df.json_str, schema).alias("data")).select("data.*") # 转换为强类型Dataset rock_dataset = parsed_df.as(rock_encoder)
4. 通过类属性访问数据
现在可以直接通过类属性访问Dataset中的数据,示例操作如下:
# 筛选并打印name属性 rock_dataset.select(RockRecord.name).show() # 遍历Dataset访问属性 for record in rock_dataset.collect(): print(f"名称: {record.name}, 岩石高度: {record.rockHeight}")
补充:Scala环境下的实现
如果是在Scala环境中,直接定义Case Class即可,Spark会自动推导Encoder,实现更简洁:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.types._ import org.apache.spark.sql.functions.from_json case class RockRecord( name: String, types: String, url: String, image: String, pullTime: String, rockHeight: String, dateLive: String, hitTime: String, desc: String ) val spark = SparkSession.builder().appName("RockDataset").getOrCreate() val schema = StructType(Seq( StructField("name", StringType, true), StructField("types", StringType, true), StructField("url", StringType, true), StructField("image", StringType, true), StructField("pullTime", StringType, true), StructField("rockHeight", StringType, true), StructField("dateLive", StringType, true), StructField("hitTime", StringType, true), StructField("desc", StringType, true) )) val rockDataset = spark.read.textFile("D:/pythonwork/FUN/input/*") .select(from_json($"value", schema).as[RockRecord]) // 访问类属性示例 rockDataset.select(RockRecord.name).show()
内容的提问来源于stack exchange,提问作者Carol
相关产品推荐
相关产品推荐

