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

Spark 2.11与1.6处理损坏JSON行的行为差异及问题

Handling Corrupt Snappy JSON Records Across Spark Versions

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 07:21:30