- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
11-05-2025 04:35 AM
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
associationsis 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:
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
associationsis 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 beTimestampType). -
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.