Location of WAL in Spark Structured Streaming

Viewed 985

I have enabled WAL for my Structured Streaming Application. Where do I find the location of WAL logs? I am able to see WAL for my Spark streaming process in the prefix receivedBlockMetadata . But, I don't see any prefix created for Structured Streaming

2 Answers

According to my understanding, WAL only works in spark streaming, not structred streaming. Structured streaming implements fault tolerance based on checkpoint like flink global state. The checkpoint stores all the state including kafka offsets and others.The location is specified in your code .

In Spark Structure Streaming, now WAL with every message from the receiver. Only two logs with metadata for every batch: offset and commit logs. You can find details of implementation in org.apache.spark.sql.execution.streaming.StreamExecution. ->

/**
   * A write-ahead-log that records the offsets that are present in each batch. In order to ensure
   * that a given batch will always consist of the same data, we write to this log *before* any
   * processing is done.  Thus, the Nth record in this log indicated data that is currently being
   * processed and the N-1th entry indicates which offsets have been durably committed to the sink.
   */
  val offsetLog = new OffsetSeqLog(sparkSession, checkpointFile("offsets"))

  /**
   * A log that records the batch ids that have completed. This is used to check if a batch was
   * fully processed, and its output was committed to the sink, hence no need to process it again.
   * This is used (for instance) during restart, to help identify which batch to run next.
   */
  val commitLog = new CommitLog(sparkSession, checkpointFile("commits"))

Both of them available in a checkpointLocation in folder offsets and commits. In Structure Streaming logs contain only offset information.

Related