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

如何在PySpark中展平嵌套JSON数据?

问题描述

我有如下格式的JSON数据:

[
 {  "student_id": 1234, 
    "room_id": "abc", 
    "enrolled": false
 }, 
 {   "student_id": 4321, 
    "room_id": "def", 
    "enrolled": true, 
    "enrollment": {
        "type": "home", 
        "date": "01-01-2020"
    }
},
 {   "student_id": 678,
    "room_id": "htf",
    "sports": {
        "team": "hockey",
        "position": "forward"
    }
}
] 

我通过以下代码实现了部分展平:

df = sc.parallelize(data).map(lambda x: json.dumps(x))

得到的结果如下:

student_idroom_idenrolledenrollmentsports
1234abcfalseNULLNULL
4321deftrue{home, 01-01-2020}NULL
678htfNULLNULL{hockey, forward}

需要进一步展平,得到如下格式的结果:

student_idroom_idenrolledtypedateteamposition
1234abcfalseNULLNULLNULLNULL
4321deftruehome01-01-2020NULLNULL
678htfNULLNULLNULLhockeyforward

解决方案

你可以利用PySpark的结构化数据操作直接展开嵌套字段,以下是两种可行的实现方式:

方式一:手动定义Schema解析JSON

如果需要严格控制数据类型,可以先定义Schema,再解析JSON字符串:

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, IntegerType, StringType, BooleanType
from pyspark.sql.functions import from_json, col

# 初始化SparkSession(若未初始化)
spark = SparkSession.builder.appName("FlattenNestedJSON").getOrCreate()

# 定义JSON结构的Schema
schema = StructType([
    StructField("student_id", IntegerType(), nullable=True),
    StructField("room_id", StringType(), nullable=True),
    StructField("enrolled", BooleanType(), nullable=True),
    StructField("enrollment", StructType([
        StructField("type", StringType(), nullable=True),
        StructField("date", StringType(), nullable=True)
    ]), nullable=True),
    StructField("sports", StructType([
        StructField("team", StringType(), nullable=True),
        StructField("position", StringType(), nullable=True)
    ]), nullable=True)
])

# 将RDD转为DataFrame并解析JSON字符串
json_rdd = sc.parallelize(data).map(lambda x: json.dumps(x))
df = json_rdd.toDF("json_str")
df = df.select(from_json(col("json_str"), schema).alias("data")).select("data.*")

# 展平嵌套字段
flattened_df = df.select(
    "student_id",
    "room_id",
    "enrolled",
    col("enrollment.type").alias("type"),
    col("enrollment.date").alias("date"),
    col("sports.team").alias("team"),
    col("sports.position").alias("position")
)

# 查看结果
flattened_df.show()

方式二:自动推断Schema(简化写法)

如果不需要严格指定数据类型,可以让Spark自动推断JSON结构,省去手动定义Schema的步骤:

from pyspark.sql import SparkSession

# 初始化SparkSession
spark = SparkSession.builder.appName("FlattenNestedJSON").getOrCreate()

# 直接将JSON数据转为结构化DataFrame
df = spark.read.json(sc.parallelize(data))

# 直接提取嵌套字段完成展平
flattened_df = df.select(
    "student_id",
    "room_id",
    "enrolled",
    "enrollment.type",
    "enrollment.date",
    "sports.team",
    "sports.position"
)

# 查看结果
flattened_df.show()

两种方式都会生成你需要的展平结果,缺失的字段会自动填充为NULL。


内容的提问来源于stack exchange,提问作者unlocknew

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 23:29:50