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

Databricks拆分大JSON文件时Spark _corrupt_record错误求助

解决Databricks中拆分大JSON文件时的Spark错误

问题描述

在Databricks运行Python代码拆分大JSON文件为两个文件时,触发以下错误:

Since Spark 2.3, the queries from raw JSON/CSV files are disallowed when the
referenced columns only include the internal corrupt record column
(named _corrupt_record by default). For example:
spark.read.schema(schema).csv(file).filter($"_corrupt_record".isNotNull).count()
and spark.read.schema(schema).csv(file).select("_corrupt_record").show().
Instead, you can cache or save the parsed results and then send the same query.
For example, val df = spark.read.schema(schema).csv(file).cache() and then
df.filter($"_corrupt_record".isNotNull).count().

已尝试在读取文件代码末尾添加.cache(),但错误依然存在。

原代码

from pyspark.sql.functions import explode, col

# Read the JSON file from Databricks storage
df_json = spark.read.json("/mnt/BigData_JSONFiles/new_test.json")
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "false")


# Convert the dataframe to a dictionary
data = df_json.toPandas().to_dict()

# Split the data into two parts
d1 = dict(itertools.islice(data.items(), 8))
d2 = dict(itertools.islice(data.items(), 8, len(data.items())))

# Convert the first part of the data back to a dataframe
df1 = spark.createDataFrame([d1])

# Write the first part of the data to a JSON file in Databricks storage
df1.write.format("json").save("/mnt/BigData_JSONFiles/new_test_header.json")

# Convert the second part of the data back to a dataframe
df2 = spark.createDataFrame([d2])

# Write the second part of the data to a JSON file in Databricks storage
df2.write.format("json").save("/mnt/BigData_JSONFiles/new_test_detail.json")

大JSON文件示例

{
  "reporting_entity_name": "launcher",
  "reporting_entity_type": "launcher",
  "plan_name": "launched",
  "plan_id_type": "hios",
  "plan_id": "1111111111",
  "plan_market_type": "individual",
  "last_updated_on": "2020-08-27",
  "version": "1.0.0",
  "in_network": [
    {
      "negotiation_arrangement": "ffs",
      "name": "Boosters",
      "billing_code_type": "CPT",
      "billing_code_type_version": "2020",
      "billing_code": "27447",
      "description": "Boosters On Demand",
      "negotiated_rates": [
        {
          "provider_groups": [
            {
              "npi": [
                0
              ],
              "tin": {
                "type": "ein",
                "value": "11-1111111"
              }
            }
          ],
          "negotiated_prices": [
            {
              "negotiated_type": "negotiated",
              "negotiated_rate": 123.45,
              "expiration_date": "2022-01-01",
              "billing_class": "organizational"
            }
          ]
        }
      ]
    }
  ]
}

解决方案

1. 错误根源

原代码的核心问题:

  • 将Spark DataFrame转成Pandas再转字典的操作,会触发Spark重复读取原始JSON文件,若解析过程生成_corrupt_record列,后续操作仅关联该列就会触发错误。
  • 按字典键数量拆分的逻辑不合理,Spark DataFrame转成的字典是键:列表结构,直接切片会导致数据结构混乱。

2. 优化代码(纯Spark操作,无Pandas转换)

直接用Spark的列选择功能拆分数据,避免转换带来的问题:

# 读取JSON文件
df_json = spark.read.json("/mnt/BigData_JSONFiles/new_test.json")
# 执行count()触发缓存(仅.cache()不会立即生效,需行动算子触发)
df_json.count()

# 定义表头字段列表
header_cols = [
    "reporting_entity_name", "reporting_entity_type", "plan_name",
    "plan_id_type", "plan_id", "plan_market_type", "last_updated_on", "version"
]
# 生成表头DataFrame并写入文件
df1 = df_json.select(header_cols)
df1.write.mode("overwrite").format("json").save("/mnt/BigData_JSONFiles/new_test_header.json")

# 生成详情DataFrame并写入文件(仅保留in_network字段)
df2 = df_json.select("in_network")
df2.write.mode("overwrite").format("json").save("/mnt/BigData_JSONFiles/new_test_detail.json")

3. 关键修复点

  • 移除Pandas转换:全程用Spark分布式操作,避免重复读取原始文件,彻底规避_corrupt_record相关错误。
  • 有效缓存:通过count()行动算子触发缓存,确保解析结果被持久化,而非仅调用.cache()。
  • 明确列选择:直接指定拆分的列,避免字典切片带来的结构错误,同时确保操作不涉及_corrupt_record列。

4. 处理JSON解析异常(可选)

若原始JSON存在格式问题导致_corrupt_record生成,建议指定Schema读取:

from pyspark.sql.types import StructType, StructField, StringType, ArrayType, DateType

# 定义嵌套结构Schema
in_network_schema = StructType([
    StructField("negotiation_arrangement", StringType()),
    StructField("name", StringType()),
    StructField("billing_code_type", StringType()),
    StructField("billing_code_type_version", StringType()),
    StructField("billing_code", StringType()),
    StructField("description", StringType())
])

# 主Schema
main_schema = StructType([
    StructField("reporting_entity_name", StringType()),
    StructField("reporting_entity_type", StringType()),
    StructField("plan_name", StringType()),
    StructField("plan_id_type", StringType()),
    StructField("plan_id", StringType()),
    StructField("plan_market_type", StringType()),
    StructField("last_updated_on", DateType()),
    StructField("version", StringType()),
    StructField("in_network", ArrayType(in_network_schema))
])

# 指定Schema读取JSON
df_json = spark.read.schema(main_schema).json("/mnt/BigData_JSONFiles/new_test.json")
df_json.count()  # 触发缓存

内容的提问来源于stack exchange,提问作者Anjaney Shrivastava

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 08:30:41