You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将类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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.24 14:06:19