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

JSON字符串转PySpark DataFrame报错,求正确转换方法

解决JSON字符串转PySpark DataFrame的问题

问题分析

你的代码存在两个关键错误导致转换失败:

  1. JSON格式语法错误:第二个JSON对象{"id":2, "name":"test2"}与第三个对象之间缺少逗号,不符合JSON数组的语法规范。
  2. RDD创建错误:sc.parallelize(some_json_string)会将字符串拆分为单个字符的RDD,PySpark无法识别这种格式的输入,最终生成_corrupt_record列。

修正后的代码

import findspark
findspark.init()
from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("MyApp") \
    .getOrCreate()

sc = spark.sparkContext
# 修正JSON字符串,补上缺失的逗号
some_json_string = """
[
  {"id":1, "name":"test1"},
  {"id":2, "name":"test2"},
  {"id":3, "name":"test3"}
]
"""

# 将JSON字符串作为单个元素传入parallelize,确保multiLine模式正确解析
df = spark.read.option("multiLine", "true").json(sc.parallelize([some_json_string]))

df.printSchema()
df.show()

另一种简化方法(无需RDD)

你也可以直接用StringIO读取JSON字符串,代码更简洁:

import findspark
findspark.init()
from pyspark.sql import SparkSession
from io import StringIO

spark = SparkSession.builder \
    .appName("MyApp") \
    .getOrCreate()

some_json_string = """
[
  {"id":1, "name":"test1"},
  {"id":2, "name":"test2"},
  {"id":3, "name":"test3"}
]
"""

df = spark.read.option("multiLine", "true").json(StringIO(some_json_string))

df.printSchema()
df.show()

正确输出结果

修正后运行代码,会得到符合预期的DataFrame:

root
 |-- id: long (nullable = true)
 |-- name: string (nullable = true)

+---+-----+
| id| name|
+---+-----+
|  1|test1|
|  2|test2|
|  3|test3|
+---+-----+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 13:22:06