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
相关产品推荐
相关产品推荐

