CertSafari
    Databricks Certified Associate Developer for Apache Spark· Lessons

    Domain 5 · Lesson 26/32

    Structured Streaming Output Sinks: Delta Tables, Kafka, Console, Memory and foreachBatch

    Create and write Streaming DataFrames and Streaming Datasets, including the basic output modes and output sinks.

    10 min read
    3.12% of exam
    5 sources
    Published 3 Oct 2026
    Docs as of 30 Sep 2026

    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.

    Streaming a groupBy count into a Delta table in complete modepython
    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")
    )

    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.

    Output mode support for the two sinks the sources document
    Sinkappendupdatecomplete
    Delta Lake / Unity Catalog managed tableSupported (default)Not supportedSupported (replaces the table each batch)
    KafkaSupportedSupportedSupported

    Checkpoint 2 of 6· Check yourself

    A streaming aggregation must write to a Unity Catalog managed table. Which output modes can it use?

    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() ) ```

    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.

    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?

    Sources34

    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.

    A foreachBatch function that skips empty micro-batchespython
    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?

    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() ```

    Sources5

    Exam traps

    Each one states something that sounds right. Open it to see what is actually true.

    1. 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. 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. 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. 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. 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. 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. 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. 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. 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. 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. 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

    Spotted a mistake, or was something unclear? Tell us.