Databricks拆分大JSON文件时Spark _corrupt_record错误求助
问题描述
在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

