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

Spark Schema校验JSON时需显式添加_corrupt_record列的原因

Understanding Spark's JSON Parsing Behavior with Corrupt Records

Great question! This behavior boils down to how Spark handles schema inference vs. explicit schema definitions, paired with its default parsing modes. Let's break it down step by step:

1. What happens during automatic schema inference?

When you let Spark infer the schema from a JSON file (by not specifying a schema), Spark does more than just guess the structure of valid records. It automatically accounts for potential malformed data by:

  • Enabling the default PERMISSIVE parsing mode (which allows handling bad records instead of failing immediately)
  • Automatically injecting the _corrupt_record column into the inferred schema. This column is of type String and stores the full raw JSON string of any record that fails to parse against the inferred structure.

Spark does this because schema inference is meant to be a flexible, "just works" experience—it anticipates that your data might have messy records and builds in this safety net by default.

2. Why explicit schemas require manual _corrupt_record setup

When you specify a custom schema, Spark treats it as an authoritative definition of your data's structure. It won't modify your schema behind the scenes, even if it encounters malformed records. Here's what happens in two scenarios:

  • Without _corrupt_record in your schema: In PERMISSIVE mode (still the default), Spark can't find a column to store the bad record content, so it sets all fields in the malformed record to null—effectively discarding the original raw data.
  • With _corrupt_record added to your schema: Now Spark has a designated place to store the raw JSON of failed records. It will set valid fields to null (since the record didn't match the schema) but preserve the original bad data in _corrupt_record.

Example Code Snippets

To make this concrete, here's how the code behaves in each case:

Automatic Schema Inference (captures bad records)

val dfInfer = spark.read.json("path/to/mixed-json-data")
dfInfer.printSchema()
// Output will include _corrupt_record: string (nullable = true) alongside your valid fields

Explicit Schema Without Corrupt Column (loses bad record data)

import org.apache.spark.sql.types.{StructType, StructField, IntegerType, StringType}

val customSchema = StructType(Seq(
  StructField("user_id", IntegerType),
  StructField("username", StringType)
))

val dfNoCorrupt = spark.read.schema(customSchema).json("path/to/mixed-json-data")
// Malformed records will have user_id: null, username: null — no trace of the original bad JSON

Explicit Schema With Corrupt Column (preserves bad records)

val customSchemaWithCorrupt = StructType(Seq(
  StructField("user_id", IntegerType),
  StructField("username", StringType),
  StructField("_corrupt_record", StringType)
))

val dfWithCorrupt = spark.read.schema(customSchemaWithCorrupt).json("path/to/mixed-json-data")
// Malformed records have user_id: null, username: null, _corrupt_record: "<raw bad JSON string>"

3. Key Configuration Note

Spark uses the spark.sql.columnNameOfCorruptRecord config to define the name of this error column (default is _corrupt_record). If you want to use a different name (like bad_record), you can set this config before reading the data, and then add that named column to your explicit schema.


内容的提问来源于stack exchange,提问作者Sanyam Jain

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:58:28