Spark Structured Streaming checkpoint恢复数据报错求助
Hey there, let's break down this checkpoint recovery issue you're hitting with Spark Structured Streaming and CSV file sources. That error message is telling you your current query setup isn't compatible with the existing checkpoint data—here's how to diagnose and fix it:
Common Root Causes & Fixes
1. Schema Mismatch (Most Likely Culprit)
If you're using inferSchema=true or changing your explicit schema between restarts, Spark can't reconcile the checkpoint's stored schema with the new one. This is a hard no for checkpoint recovery.
Fix:
Always use an explicit, fixed schema and never rely on schema inference. For example:
import org.apache.spark.sql.types._ val schema = StructType(Seq( StructField("order_id", IntegerType, nullable = false), StructField("customer_name", StringType, nullable = true), StructField("order_amount", DoubleType, nullable = false), StructField("order_timestamp", TimestampType, nullable = false) ))
Ensure this exact schema is used every time you start the stream—no tweaks to field names, types, or nullability allowed.
2. Changed Stream Configuration
Any modification to your readStream options between restarts will break checkpoint compatibility. This includes:
- Altering
maxFilesPerTriggervalue - Changing CSV-specific options like
header,sep,quote, orescape - Switching output modes (e.g., from
appendtocomplete)
Fix:
Double-check that every .option() in your readStream and writeStream calls matches exactly what you used when you first created the checkpoint. Even a tiny change (like removing header=true) will trigger this error.
3. Corrupted Checkpoint Data
If your initial stream crashed unexpectedly or failed to write the checkpoint properly, the stored offset data might be corrupted.
Fix:
As the error suggests, delete the offsets subfolder in your checkpoint directory first:
# Windows command to delete the folder rmdir /s /q C:/Users/q794089/Documents/Hadoop/SparkScala/recoveringcheckpoint/checkpoint/offsets
If that doesn't work, delete the entire checkpoint folder and restart fresh (just make sure your schema and config are locked in first).
Example of a Compatible Stream Setup
Here's how your code should look to ensure checkpoint recovery works reliably:
import org.apache.spark.sql.SparkSession object OrderStreamRecovery { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("OrderStreamWithCheckpoint") .master("local[*]") .getOrCreate() // Fixed, explicit schema (no changes allowed between restarts) val schema = StructType(Seq( StructField("order_id", IntegerType, nullable = false), StructField("customer_id", StringType, nullable = true), StructField("amount", DoubleType, nullable = false), StructField("order_date", TimestampType, nullable = false) )) val orderStream = spark.readStream .format("csv") .schema(schema) .option("header", "true") .option("sep", ",") .option("maxFilesPerTrigger", 2) // Fixed value—don't change later .load("path/to/your/csv/folder") val query = orderStream.writeStream .format("console") .outputMode("append") .option("checkpointLocation", "C:/Users/q794089/Documents/Hadoop/SparkScala/recoveringcheckpoint/checkpoint") .start() query.awaitTermination() } }
Key Takeaway
Spark Structured Streaming's checkpoint recovery for file sources relies on absolute consistency between your initial stream configuration and every restart. Lock in your schema and all stream options, and you'll avoid this error in the future.
内容的提问来源于stack exchange,提问作者Nagendra Palla

