Spark 2.11与1.6处理损坏JSON行的行为差异及问题
I’ve dealt with this exact scenario before—let’s walk through how it works in Spark 1.6 vs Spark 2.x (note: Spark doesn’t have a "2.11" version; you’re likely referring to Spark 2.x running on Scala 2.11):
Reading Snappy-compressed JSON Files
First, here’s the standard code to read a Snappy-compressed JSON file using SQLContext (for Spark 1.x):
val sqlContext = new org.apache.spark.sql.SQLContext(sc) val df = sqlContext.read.json("s3://bucket/problemfile.snappy")
Spark 1.6: Automatic Corrupt Record Capture
In Spark 1.6, the JSON reader is permissive by default. It automatically captures invalid JSON records in a special _corrupt_record column, which lets you split valid and invalid data with simple queries:
// Isolate corrupted records val invalidJSON = rawEvents.select("*").where("_corrupt_record is not null") // Keep only valid records val validJSON = rawEvents.select("*").where("_corrupt_record is null")
Spark 2.x: Explicit Configuration Required
The behavior changed in Spark 2.x—by default, if the reader encounters severely malformed JSON, it might fail outright instead of populating the _corrupt_record column. That’s why you’re hitting exceptions when trying to query that column.
To restore the Spark 1.6-style behavior (capturing corrupt records instead of failing), you need to explicitly configure the JSON reader:
// Using SparkSession (the standard entry point in Spark 2.x) import org.apache.spark.sql.SparkSession val spark = SparkSession.builder() .appName("CorruptRecordHandler") .getOrCreate() val df = spark.read .option("mode", "PERMISSIVE") // This is default, but explicit is safer .option("columnNameOfCorruptRecord", "_corrupt_record") // Ensure the column is named correctly .json("s3://bucket/problemfile.snappy") // Now you can split records just like in Spark 1.6 val invalidJSON = df.select("*").where("_corrupt_record is not null") val validJSON = df.select("*").where("_corrupt_record is null")
If you’re using a custom schema, make sure to add a _corrupt_record: StringType field to your schema definition—this ensures the reader knows where to store invalid records.
内容的提问来源于stack exchange,提问作者scrayon

