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

如何解决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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 22:05:32