mark_ott
Databricks Employee
Databricks Employee

To resolve the DLT streaming aggregation error about unsupported output modes and watermarks in Databricks, you need to carefully set watermarks on the original event timestamp rather than on computed columns like "time_window" and carefully consider join semantics between streaming and static datasets. Here is a breakdown of key points and specific steps to correct your code:

Key Concepts

  • Watermark must be set on a column that directly represents event time, usually "timestamp". Setting it on a derived column like "time_window" is not supported for streaming aggregation or output mode resolution.

  • Streaming joins: You can only join a streaming DataFrame with a static one in append mode. Joining two streaming DataFrames with aggregation usually requires watermarks on both sides and only on event columns. Complete output mode is not possible for joins, only append.

  • View vs Table: Even if associations is a VIEW, as long as its source is static (not itself a stream), it's fine for joining.

Corrected Code Steps

Here's an approach that tracks Databricks recommendations:

python
import pyspark.sql.functions as F from pyspark.sql import DataFrame @Dlt.table def temp_and_humidity(): # Load streams on event tables temp = spark.readStream.table("LIVE.temp") humidity = spark.readStream.table("LIVE.humidity") associations = spark.read.table("LIVE.sensor_associations") # static view/table # Ensure schemas are compatible before union temp = temp.withColumn("humidity", F.lit(None)) humidity = humidity.withColumn("temp", F.lit(None)).withColumnRenamed("sensor_id", "sensor_group_id") # Add watermark ON THE EVENT TIME COLUMN before any aggregation temp = temp.withWatermark("timestamp", "2 minutes") # use a reasonable delay humidity = humidity.withWatermark("timestamp", "2 minutes") # Join temp stream with associations (static) temp = temp.join(associations, temp["sensor_id"] == associations["temp_sensor_id"], "left") # Union the streams temp_and_humidity = temp.unionByName(humidity) # Group by time window and sensor group # Window is derived, but watermark is set on 'timestamp' return ( temp_and_humidity .withColumn("time_window", F.window("timestamp", "1 minute")) .groupBy("time_window", "sensor_group_id") .agg( F.first("humidity", ignorenulls=True), F.first("temp", ignorenulls=True) ) )

Critical Adjustments

  • Watermark must be set on "timestamp" before aggregation, not on “time_window”.

  • Join static with stream only: When joining, ensure at least one side is static (your associations is OK).

  • Avoid joining two streams before aggregation: If you need to join the results of stream aggregations, consider changing your logic or introducing a static reference.

Troubleshooting Checklist

  • Inspect "timestamp" columns for correct type (should be TimestampType).

  • Ensure watermarks are set before aggregation (groupBy).

  • Only join streams with static tables in append mode.

Reference

For more details and troubleshooting, Databricks provides an official guide on watermarks and aggregation in streaming. Double-check their advice and make sure you are not computing watermarks on windows.


With these changes, your pipeline should conform to Databricks' streaming requirements and the exceptions about unsupported output mode should be resolved.

View solution in original post