What you will be able to do
- Write a streaming DataFrame to a Delta table with toTable() in append or complete mode
- Identify which output modes the Delta Lake and Kafka sinks accept
- Name the caveats of the console and memory sinks
- Use foreachBatch to write to a sink that has no built-in streaming support, and state its delivery guarantee
1.The Delta table sink with toTable()
On Databricks the most common sink is a Delta table, and Unity Catalog managed tables are Delta tables. You write to one with writeStream...toTable(name), or with start() and a path option for a location-based table. Delta's transaction log gives exactly-once processing even when other streams or batch jobs write to the same table. With no output mode set, the stream appends: each micro-batch adds only new rows.
Complete mode works differently with Delta: it replaces the whole table after every batch. That suits a small summary table that is always recomputed, such as event counts per customer. The example below is runnable. It aggregates a streaming read of an events table and writes the counts in complete mode. Because it uses an availableNow trigger, it processes what is there and then stops.
q = (spark.readStream
.table("main.streaming_examples.events")
.groupBy("customerId")
.count()
.writeStream
.outputMode("complete")
.option("checkpointLocation", f"/tmp/delta/eventsByCustomer/_checkpoints/{uuid.uuid4()}")
.trigger(availableNow=True)
.toTable("main.streaming_examples.events_by_customer")
)Checkpoint 1 of 6· Fill the gap
Which option name completes this append-mode write to a Delta table?
(events.writeStream
.outputMode("append")
.option(" ? ", "/tmp/delta/events/_checkpoints/")
.toTable("events")
)The checkpoint directory is set with the checkpointLocation option. path sets where the output data goes, and schemaLocation is an Auto Loader read option.
“By default, streams run in append mode and only add new records to the table.”Source: docs.databricks.com
Sources1
2.Which output modes each sink accepts
The query fails. Not every sink supports every output mode, and the Delta sink is the one to remember. Delta Lake, and therefore every Unity Catalog managed table, accepts append and complete but not update. The Databricks output-mode page shows an update-mode toTable() example and then warns that this example fails against a managed table. Kafka, by contrast, accepts all three modes. To get update-like results into Delta, Databricks points you to a merge inside foreachBatch, covered in the last section of this page.
| Sink | append | update | complete |
|---|---|---|---|
| Delta Lake / Unity Catalog managed table | Supported (default) | Not supported | Supported (replaces the table each batch) |
| Kafka | Supported | Supported | Supported |
Checkpoint 2 of 6· Check yourself
A streaming aggregation must write to a Unity Catalog managed table. Which output modes can it use?
Managed tables are Delta tables, and the Delta sink rejects update mode.
“Delta Lake sinks support only append and complete modes, so the update output mode example fails against a managed table.”Source: docs.databricks.com
Checkpoint 3 of 6· Exam question
The following streaming write fails immediately with `AnalysisException: checkpointLocation must be specified`. Which single change lets the query start successfully while keeping Parquet as the output format? ```python query = ( filtered_events.writeStream .format("parquet") .option("path", "/mnt/output/events") .outputMode("append") .start() ) ```
Correct answer: A — Add `.option("checkpointLocation", "/mnt/checkpoints/events")` before `.start()` so Spark can persist offsets and state for recovery.
- A. Fault-tolerant sinks such as the file sink require a `checkpointLocation` so Spark can record processed offsets and any intermediate state needed to resume correctly after a restart. Adding this option is exactly what the error message asks for.
- B. `append` is already the default output mode, so removing it changes nothing about the query's behavior, and the checkpoint requirement applies regardless of which output mode is configured on a fault-tolerant sink.
- C. Checkpointing is a requirement of fault-tolerant sinks in general, not something unique to the Delta format; Parquet file sinks also need a checkpoint location, so switching formats does not remove the requirement.
- D. A one-time trigger still writes through a fault-tolerant sink and therefore still needs a checkpoint location to track offsets and guarantee exactly-once output; triggers do not exempt a query from checkpointing.
- E. `DataStreamWriter` does not expose a `.save()` method after `.start()`, and the destination `path` option is unrelated to checkpoint metadata, so this does not address the missing checkpoint location.
Sources2
3.Console, memory, Kafka and file-path sinks
You pick a sink with format() and finish the write with start(). The rate-to-console sample in the DataStreamWriter reference uses df.writeStream.format("console").start(), then stops the query with q.stop() a few seconds later. That makes it a quick way to check that a stream is producing rows. The memory sink and notebook display() are similar development tools. Both quietly create a temporary checkpoint location if you leave out checkpointLocation, and that temporary checkpoint guarantees nothing.
The automatically generated temporary checkpoint does not guarantee fault tolerance or data consistency, and it might not get cleaned up properly. An explicit location fixes both problems.
For Kafka, set format("kafka"), the kafka.bootstrap.servers and topic options, and a checkpoint, then call start(). Databricks' tutorial uses this pattern to enrich Kafka data and write it back to Kafka for the lowest latency. To write files rather than a named table, pass a path option and call start(), as in the earlier tutorial example. Once a query has a checkpoint, you cannot always change its sink type freely. A file sink can later be switched to Kafka, but a Kafka sink cannot be switched to a file sink.
Checkpoint 4 of 6· Check yourself
A developer writes a stream to the memory sink and omits checkpointLocation. What do the docs say happens?
The memory sink and display() create a temporary checkpoint automatically, but that location does not ensure fault tolerance.
“These temporary checkpoint locations do not ensure any fault tolerance or data consistency guarantees”Source: docs.databricks.com
4.foreachBatch: writing to any sink with batch code
Some targets have no streaming sink at all, and some operations, such as a Delta MERGE, are not supported on streaming DataFrames. foreachBatch covers both cases. You pass it a function, and the engine calls it once per micro-batch with two arguments: a normal batch DataFrame holding that micro-batch's output, and the micro-batch's unique ID. Inside the function, you can use any batch writer. This is also how you write streaming aggregations into Delta with update-like behaviour: MERGE INTO inside foreachBatch. Databricks says you must use foreachBatch for Delta merges in streaming.
def process_batch(output_df, batch_id):
# Process valid DataFrames only
if not output_df.isEmpty():
# business logic
pass
streamingDF.writeStream.foreachBatch(process_batch).start()Three rules come with this flexibility. First, the function can receive an empty DataFrame, for example after an OPTIMIZE on a Delta source that had no files to process, so check for it as the sample does. Second, foreachBatch guarantees only at-least-once writes. To get exactly-once, you use the batch_id to deduplicate your own writes. Third, it relies on micro-batch execution, so continuous processing mode needs foreach() instead. When a query includes a stateful operator, the function must consume the entire batch DataFrame. Finally, if you need to send the same stream to several sinks, Databricks recommends separate streaming writers rather than one foreachBatch that writes them one after another.
Checkpoint 5 of 6· Check yourself
What write guarantee does foreachBatch provide on its own?
The function may run again for the same batch, so the guarantee is at-least-once unless your code uses the batch ID to make its writes idempotent.
“foreachBatch() provides only at-least-once write guarantees.”Source: docs.databricks.com
Checkpoint 6 of 6· Exam question
A team runs this streaming aggregation and then queries the in-memory result table from a notebook cell. What happens? ```python counts = clicks_stream.groupBy("page").count() query = ( counts.writeStream .queryName("page_counts") .outputMode("complete") .format("memory") .start() ) spark.sql("SELECT * FROM page_counts").show() ```
Correct answer: A — The SQL query returns the full, continuously updated count per page because complete mode rewrites the entire in-memory table after every micro-batch.
- A. The memory sink registers an in-memory table named after `queryName`, and in `complete` mode Spark rewrites that table with the full Result Table after each trigger, so querying it with `spark.sql` returns the current full aggregation. This is the documented behavior for debugging aggregations with the memory sink.
- B. The memory sink supports both `append` and `complete` output modes, and `complete` is the mode typically paired with aggregation queries on this sink, so it does not raise an exception here.
- C. `queryName` only names the in-memory table for querying; it does not override the configured output mode, and the query as written runs in `complete` mode, not append-only.
- D. The in-memory table becomes queryable as soon as the streaming query has processed at least one trigger; `spark.sql` can read it directly, and calling `awaitTermination()` is not required and would block the notebook cell instead.
- E. The memory sink is explicitly documented as not fault-tolerant and does not require a `checkpointLocation` to populate results; it keeps the full table in the driver's memory rather than relying on checkpoint state to return rows.
Sources5
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.outputMode("update") works with toTable() on a Unity Catalog managed table.Why is that wrong?
Managed tables are Delta, and the Delta sink supports only append and complete. To get update-like results, use MERGE inside foreachBatch.
Covered in Which output modes each sink accepts
2.foreachBatch inherits Structured Streaming's exactly-once guarantee automatically.Why is that wrong?
foreachBatch alone is at-least-once. Exactly-once requires deduplicating your writes with the batchId passed to the function.
Covered in foreachBatch: writing to any sink with batch code
Practise it for real
Stream an aggregation into a Delta table in complete mode and confirm the result
1.Run spark.sql("CREATE SCHEMA IF NOT EXISTS main.streaming_examples") and seed main.streaming_examples.events with customerId values c1, c1, c1, c2, c2, c3 using saveAsTable, as in the Delta Lake streaming doc.
Why: A Delta table can act as a streaming source through spark.readStream.table().
You should see: An events table with six rows.
2.Start the query from this page: readStream.table on events, groupBy("customerId").count(), writeStream with outputMode("complete"), a unique checkpointLocation, trigger(availableNow=True), and toTable("main.streaming_examples.events_by_customer").
Why: Complete mode is valid here because the query is a streaming aggregation, and the Delta sink accepts complete mode.
You should see: A StreamingQuery object is returned.
3.Call q.awaitTermination().
Why: availableNow processes all unprocessed records and then shuts down.
You should see: The call returns once the backlog is processed.
4.Run SELECT * FROM main.streaming_examples.events_by_customer ORDER BY customerId.
Why: Confirms that complete mode wrote the full aggregate result to the table.
You should see: c1 -> 3, c2 -> 2, c3 -> 1.
Stuck? Get a nudge
Try changing outputMode to "update" with a fresh checkpoint and see the query fail against the Delta sink.
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“By default, streams run in append mode and only add new records to the table.”
↩︎ The Delta table sink with toTable()“Use Structured Streaming with complete mode to replace the entire table after every batch.”
↩︎ The Delta table sink with toTable()“The Delta Lake transaction log guarantees exactly-once processing, even when there are other streams or batch queries running concurrently against the table.”
↩︎ The Delta table sink with toTable() - 2.
“Kafka supports all output modes.”
↩︎ Which output modes each sink accepts“Delta Lake sinks support only append and complete modes, so the update output mode example fails against a managed table.”
↩︎ Which output modes each sink accepts“Delta Lake, which backs all Unity Catalog managed tables, supports append and complete modes but not update mode.”
↩︎ Exam trap 1 - 3.
“the memory sink, automatically generate a temporary checkpoint location if you omit this option.”
↩︎ Console, memory, Kafka and file-path sinks“Kafka sink to file sink is not allowed.”
↩︎ Console, memory, Kafka and file-path sinks“These temporary checkpoint locations do not ensure any fault tolerance or data consistency guarantees”
↩︎ Checkpoint - 4.
“You can use Databricks to apply transformations to data ingested from Kafka and then write data back to Kafka.”
↩︎ Console, memory, Kafka and file-path sinks - 5.
“You must use foreachBatch for Delta Lake merge operations in Structured Streaming.”
↩︎ foreachBatch: writing to any sink with batch code“foreachBatch() might receive an empty DataFrame, and your code must handle this scenario.”
↩︎ foreachBatch: writing to any sink with batch code“If you write data in continuous mode, use foreach() instead.”
↩︎ foreachBatch: writing to any sink with batch code“Databricks recommends using multiple Structured Streaming writers for best parallelization and throughput.”
↩︎ foreachBatch: writing to any sink with batch code“you can use the batchId provided to the function as way to deduplicate the output and get an exactly-once guarantee”
↩︎ Exam trap 2“foreachBatch() provides only at-least-once write guarantees.”
↩︎ Checkpoint