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

Spark读取JSON列表时出现_corrupt_record错误的解决求助

问题:Spark读取JSON RDD时出现_corrupt_record错误

我通过循环从网页读取JSON数据并追加到total_results列表,将其转为RDD后用Spark读取时,持续抛出_corrupt_record错误,查阅多篇技术帖后仍未解决,求可行修复方案。

相关代码

total_results = []

response = requests.get(getURL, headers=headers)
data = response.json()
total_results.append(data) 
.......
.......
rdd = spark.sparkContext.parallelize((total_results) )
print(rdd)
df = spark.read.option('multiline','true').json(rdd)
df.show()

错误信息

ParallelCollectionRDD[189] at parallelize at PythonRDD.scala:195
                                                   
+--------------------+
|     _corrupt_record|
+--------------------+
|[{'API_UWI': '42-...|
|[{'API_UWI': '33-...|
+--------------------+
Sample data  from output

{"_corrupt_record":"[{'abc_pqq': '12-00-45672', 'abc_pqq_12': '12-00- 
45672-00', 'Unformatted': '0421733644800',............
{"_corrupt_record":"[{'abc_pqq': '13-10-45672', 'abc_pqq_12': '322-173- 
36499-00', 'Unformatted': '222223644800',..........
{"_corrupt_record":"[{'abc_pqq': '22-223-45678', 'abc_pqq_12': '22-111- 
9876543', 'Unformatted': '567890000',...................
{"_corrupt_record":"[{'abc_pqq': '33-22-678900', 'abc_pqq_12': '99-88- 
7654321', 'Unformatted': '111111111',...............
....................
................

total_results列表样本

[[{'abc_pqq': '11-111-1111', 'abc_pqq_12': '11-111-1111-1111', 
'Unformatted': '421733878600', 'abc_pqq_14': '22-222-222-22222', 
.............................................'ID': 82346790000}, 
{'abc_pqq': '11-222-2222', 'abc_pqq_12': '22-222-222-22222', 
'Unformatted': '420230106900', 'abc_pqq_14': '44-444-444-444444', 
'..............................................'

解决建议

1. 扁平化嵌套列表

从样本可看出,total_results是嵌套列表(外层列表的每个元素又是包含多条数据的子列表),直接并行化后Spark会把整个子列表当作单个JSON对象解析,导致失败。修改循环中的追加逻辑:

total_results = []
for ...:  # 你的循环逻辑
    response = requests.get(getURL, headers=headers)
    data = response.json()
    # 用extend替代append,将子列表的元素逐一加入total_results
    total_results.extend(data)

2. 生成合法JSON字符串RDD

Spark的json()读取器要求RDD每个元素是标准JSON字符串,而非Python字典。需将每个字典转为JSON字符串:

import json

# 扁平化后,把每个字典转为双引号格式的JSON字符串
rdd = spark.sparkContext.parallelize([json.dumps(item) for item in total_results])
df = spark.read.json(rdd)
df.show()

3. 简化方案:直接创建DataFrame

如果数据量不大,可跳过RDD步骤,直接用字典列表创建DataFrame:

# 扁平化后的total_results是字典列表,直接生成DataFrame
df = spark.createDataFrame(total_results)
df.show()

4. 校验JSON合法性

若仍报错,检查每个JSON字符串:

  • 确保所有引号为双引号(JSON语法要求)
  • 无语法错误(如缺失逗号、括号不闭合)

内容的提问来源于stack exchange,提问作者Arun.K

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 15:45:23