Interview/Spark

Spark interview questions: Structured Streaming: an unbounded table

Interactive Spark interview questions on Structured Streaming: an unbounded table. Practice with the matching lesson. A stream is an unbounded table read in micro-batches. Checkpoints give exactly-once; watermarks bound state.

Lesson · Simulation

Your streaming query looks like ordinary DataFrame code. What is the engine actually doing under it?

Answer it out loud, then reveal. Play steps through like the simulators.

Production scenario

Structured Streaming: an unbounded table

Streaming aggregation OOMs after four days of clean running

Symptoms

  • Runs fine for days, then executors die and restart in a loop
  • Batch duration climbing steadily in the days before the failure
  • State row count grows every batch and never comes down

All questions on this page

Indexed as FAQ. Open any item if you prefer a list to Play.

beginner

Your streaming query looks like ordinary DataFrame code. What is the engine actually doing under it?

Structured Streaming: an unbounded table · tap to open the answer

Short: Running a series of micro-batches: each batch reads a new offset range and executes the same incremental plan.

Detailed: The engine records which offsets it has processed in the checkpoint's offset log, plans a batch over the new range, and writes a commit when the batch finishes. That is why Spark UI shows a steady stream of small jobs rather than one long-running job.

Common mistake: Describing it as row-by-row streaming, like a message-queue consumer loop.

Follow-up: What guarantees the same row is not processed twice after a restart?

Lesson · Simulation

intermediate

On-call deleted the checkpoint directory to 'clear a stuck stream'. What did they just do?

Structured Streaming: an unbounded table · tap to open the answer

Short: Threw away the offset log, the commit log, and all state — the stream now restarts with no memory of what it processed.

Detailed: The checkpoint location holds offsets, commits, and the state store, and exactly-once depends on those lining up with the sink. Deleting it means reprocessing from the configured starting position — duplicates into an append sink, or silently skipped data if it starts at latest — and every aggregation's state is gone.

Common mistake: Treating the checkpoint directory as a cache you can safely clear when something looks stuck.

Follow-up: When is deleting a checkpoint actually the correct move?

Lesson · Simulation

senior

A stateful aggregation ran clean for a week, then started OOMing every few hours. First hypothesis?

Structured Streaming: an unbounded table · tap to open the answer

Short: No watermark, so state retained every key forever and finally outgrew the executors.

Detailed: Without withWatermark the engine cannot know a group will never receive late data, so its state is kept indefinitely. The streaming progress metrics show the state operator's numRowsTotal climbing monotonically and state store commit time growing with it. A watermark lets keys older than the event-time threshold be dropped.

Senior: Events later than the watermark are dropped and cannot be retrofitted — you are signing up to a bounded lateness contract, so route stragglers to a repair batch job if the business needs them.

Common mistake: Restarting the stream on a bigger cluster, which just moves the failure a few days out.

Follow-up: What do you give up when you add the watermark?

Lesson · Simulation

architect

Your foreachBatch writes to a transactional store and a restart happens mid-batch. How do you keep the sink correct?

Structured Streaming: an unbounded table · tap to open the answer

Short: Make the write idempotent and keyed, because the same batchId can be re-executed.

Detailed: foreachBatch hands you the DataFrame and the batchId, and on failure the engine may replay a batch that already wrote part of its rows. Use business keys (or the batchId) in a MERGE / upsert, or record the last committed batchId in the sink transactionally so a replay becomes a no-op. Output mode matters too: append emits only new rows, update emits changed keys, complete rewrites the whole result and only works for small aggregates.

Senior: When latency tolerance is hours: AvailableNow consumes everything currently available across several micro-batches and stops, so you keep streaming semantics and checkpointed offsets without paying for a cluster that never sleeps.

Common mistake: Assuming the checkpoint's exactly-once guarantee extends to an arbitrary external sink you wrote yourself.

Follow-up: When would you pick Trigger.AvailableNow over an always-on stream?

Lesson · Simulation