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

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 maxFilesPerTrigger value
  • Changing CSV-specific options like header, sep, quote, or escape
  • Switching output modes (e.g., from append to complete)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 09:10:37