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

如何从JSON文件的user节点创建Spark DataFrame?含优化方案

问题解答

一、正确传入数据到createDataFrame的方法

首先你的JSON文件存在语法错误,必须先修正:

  • docID行末尾缺少逗号
  • user数组未闭合,整个JSON末尾缺少}

修正后的JSON示例:

{
     "docID":"dfter-yutr",
     "user" : [ [ "2111", "OLIVIA", "37" ],
              [ "2112", "AVA", "47" ],
              [ "2113", "ISABELLA", "55" ] ]
}

另外,你定义的schema = ["id","Name", "Age"]仅指定列名,无法约束数据类型,推荐用Spark的StructType定义带类型的Schema,避免后续类型转换问题:

from pyspark.sql.types import StructType, StructField, StringType, IntegerType

# 定义用户数据的Schema,将Age指定为IntegerType而非默认字符串
user_schema = StructType([
    StructField("id", StringType(), nullable=False),
    StructField("Name", StringType(), nullable=False),
    StructField("Age", IntegerType(), nullable=False)
])

如果要用createDataFrame方法,需先提取user节点的数组数据:

# 读取原始JSON文件
raw_df = spark.read.json("path/to/your/json/file.json")

# 提取user列的数组数据(仅适合小数据集,大数据集会占用Driver内存)
user_data_list = raw_df.select("user").first()[0]

# 创建DataFrame
df = spark.createDataFrame(user_data_list, schema=user_schema)

⚠️ 注意:这种方式会把整个user数组加载到Driver节点内存,仅适合小数据集,不适合10万条记录的场景。


二、10万条记录的最优DataFrame创建方案

针对大数据量,必须采用分布式处理方案,避免Driver内存溢出,步骤如下:

1. 定义完整的嵌套Schema

提前指定顶层JSON和嵌套user数组的Schema,避免Spark自动推断类型(自动推断会扫描全量数据,耗时且易出错):

from pyspark.sql.types import StructType, StructField, StringType, ArrayType, IntegerType

# 定义user数组中单个元素的Schema
user_item_schema = StructType([
    StructField("id", StringType(), nullable=False),
    StructField("Name", StringType(), nullable=False),
    StructField("Age", IntegerType(), nullable=False)
])

# 定义整个JSON文件的顶层Schema
top_level_schema = StructType([
    StructField("docID", StringType(), nullable=False),
    StructField("user", ArrayType(user_item_schema), nullable=False)
])

2. 分布式读取并展开数组

使用Spark原生的explode函数展开嵌套数组,全程分布式执行,不占用Driver内存:

from pyspark.sql.functions import explode, col

# 读取JSON文件,指定提前定义的Schema
raw_df = spark.read.schema(top_level_schema).json("path/to/your/json/file.json")

# 展开user数组,将每个子数组转为单独的行,再提取列
user_df = raw_df.select(explode(col("user")).alias("user_record")) \
                .select("user_record.*")

方案优势

  • 分布式处理:所有操作在Executor节点并行执行,避免Driver内存瓶颈
  • Schema预定义:减少Spark类型推断的开销,确保数据类型准确(如Age转为整数)
  • 高效可靠:原生Spark算子性能优异,适合处理10万条及更大规模的数据

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 19:40:03