Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Structured Streaming output modes

PySpark · Streaming & Newer Spark

Structured Streaming output modes

Hardpyspark-84
structured-streamingoutput-modeswatermark

Question

Explain append, update and complete output modes.

Solution

The output mode says which rows of the result table Spark writes to the sink at each trigger. Append writes only new, final rows. Update writes rows that changed. Complete rewrites the whole result every time.

Setting it

(agg_stream.writeStream
    .outputMode("update")
    .format("console")
    .start())

What each mode sends

Take a running count of orders per region. In micro-batch 1 the counts are north=2, south=1. In micro-batch 2, two more north orders arrive, and one east order.

Mode       Batch 2 output
append     (not allowed for this query without a watermark)
update     north=4, east=1        (only rows that changed)
complete   north=4, south=1, east=1   (the full result table)

Append

A row is written once, when Spark knows it will never change. For plain row-by-row transformations (filter, map, select), every row is final at once, so append is natural. For aggregations over time windows, a row is only final after the watermark passes the end of the window. So append with an aggregation needs a watermark, and the results appear late, after the window closes. Without a watermark, Spark rejects the query, since it can never be sure a count is done.

Update

Writes only rows that were new or changed in this trigger. Good for a dashboard table or a key-value store that you upsert into. With no aggregation, it behaves like append.

Complete

Rewrites everything each trigger. It requires an aggregation, and Spark must keep the full result in state, so it only suits small results such as a count per country.

Sink support

Not every sink supports every mode. The file sink (Parquet, JSON and so on) supports only append, because files cannot be updated in place. Sinks such as Delta (through foreachBatch and MERGE) or a database give you update behaviour. Console and memory sinks support all three, which is why tutorials use them.

What to say

Choose by the question "can a row change after I write it". If no, append (with a watermark for aggregations). If yes and the sink can upsert, update. Complete only for small summaries.

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext