如何解决PySpark读取API多行JSON返回_corrupted_record的问题
解决PySpark存储嵌套/多行API JSON时出现
_corrupted_record的问题 问题描述
使用PySpark存储API返回的JSON文件时,无嵌套结构的响应能正常处理,但遇到嵌套或多行JSON时,生成的DataFrame会出现_corrupted_record字段,解析失败。当前代码及错误输出如下:
当前代码
response = requests.request("POST", url, headers=headers, data=payload) if response.status_code == 200 and len(response.json()) != 0: file = spark.read.json(sc.parallelize([response.json()])) file.write.mode('overwrite').option("multiline","true").format("json").save(dumpPath) file.show() else: if response.status_code != 200: print(f"Failed to get data from API. Status code: {response.status_code}") else: print(f"Delta is Empty. Status code: {response.status_code}")
错误输出
+--------------------+ | _corrupt_record| +--------------------+ |[{'encodedKey': '...| +--------------------+
原因分析
问题出在spark.read.json的使用逻辑上:
- 你先通过
response.json()把API响应转换成了Python字典/列表对象,再用sc.parallelize([response.json()])生成RDD。 - 但
spark.read.json的设计目标是解析JSON格式的字符串,而非Python原生数据结构。当遇到嵌套结构或数组时,Spark会把整个Python对象当作单个字符串来解析,自然会出现解析错误,生成_corrupted_record。
解决方案
提供两种可靠的修复方式,按需选择:
方式一:直接使用API返回的原始JSON文本
跳过Python对象转换,直接用response.text获取原始JSON字符串,让Spark原生解析器处理嵌套结构,同时开启multiline选项支持多行JSON:
response = requests.request("POST", url, headers=headers, data=payload) if response.status_code == 200 and response.text.strip(): # 用原始JSON文本创建RDD json_rdd = sc.parallelize([response.text]) # 读取时开启multiline,确保嵌套/多行JSON能被正确解析 file = spark.read.option("multiline", "true").json(json_rdd) file.write.mode('overwrite').option("multiline","true").format("json").save(dumpPath) file.show() else: if response.status_code != 200: print(f"Failed to get data from API. Status code: {response.status_code}") else: print(f"Delta is Empty. Status code: {response.status_code}")
方式二:用Python对象直接创建DataFrame
既然已经把JSON转成了Python对象,直接用spark.createDataFrame创建DataFrame,这种方式能直接识别Python字典/列表的嵌套结构:
response = requests.request("POST", url, headers=headers, data=payload) if response.status_code == 200 and len(response.json()) != 0: data = response.json() # 如果API返回的是单个字典(而非数组),需要包装成列表 if isinstance(data, dict): data = [data] # 直接用Python对象列表创建DataFrame file = spark.createDataFrame(data) file.write.mode('overwrite').option("multiline","true").format("json").save(dumpPath) file.show() else: if response.status_code != 200: print(f"Failed to get data from API. Status code: {response.status_code}") else: print(f"Delta is Empty. Status code: {response.status_code}")
注意事项
- 方式一更适合处理大体积、结构复杂的JSON,Spark的原生解析器对嵌套结构的支持更稳定。
- 方式二更灵活,如果你需要在创建DataFrame前对数据做预处理(比如字段过滤、类型转换),可以优先选择这种方式。
内容的提问来源于stack exchange,提问作者PedroCSeq
相关产品推荐
相关产品推荐

