JSON字符串转PySpark DataFrame报错,求正确转换方法
解决JSON字符串转PySpark DataFrame的问题
问题分析
你的代码存在两个关键错误导致转换失败:
- JSON格式语法错误:第二个JSON对象
{"id":2, "name":"test2"}与第三个对象之间缺少逗号,不符合JSON数组的语法规范。 - 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
相关产品推荐
相关产品推荐

