You are developing a Structured Streaming job in Databricks that reads JSON data from an Auto Loader source, transforms the schema, and writes the output to a Delta Lake table. During execution, downstream consumers complain that the data contains duplicates due to upstream retries. Which operation should you apply to the DataFrame to ensure exactly-once processing semantics before writing to the Delta table?
Trap 1: Apply the distinct() transformation on the entire streaming…
Calling distinct() on an unbounded streaming DataFrame requires Spark to maintain an infinite state of all previously seen records to check for duplicates, which will eventually cause the streaming job to fail due to out-of-memory errors.
Trap 2: Use the dropDuplicates() method with a set of unique identifier…
Unbounded dropDuplicates() operations accumulate state indefinitely as new data arrives, mirroring the memory growth issues of distinct(), making it unsuitable for long-running production streaming applications processing continuous streams of data.
Trap 3: Enable Delta Lake change data feed and configure the sink to…
Delta Lake change data feed tracks row-level changes for downstream consumers rather than automatically deduplicating incoming streaming micro-batches, meaning duplicate records from upstream sources will still be persisted unless handled in the Spark logic.
- A
Apply the distinct() transformation on the entire streaming DataFrame prior to writing the output to the Delta table.
Why it fails: Calling distinct() on an unbounded streaming DataFrame requires Spark to maintain an infinite state of all previously seen records to check for duplicates, which will eventually cause the streaming job to fail due to out-of-memory errors.
- B
Use the dropDuplicates() method with a set of unique identifier columns without specifying a watermark.
Why it fails: Unbounded dropDuplicates() operations accumulate state indefinitely as new data arrives, mirroring the memory growth issues of distinct(), making it unsuitable for long-running production streaming applications processing continuous streams of data.
- C
Configure a watermark on the event-time column and apply dropDuplicatesWithinWatermark() using the unique identifier columns.
Using dropDuplicatesWithinWatermark() alongside a watermark allows Spark to safely discard old state data outside the window threshold, preventing memory exhaustion while successfully filtering out duplicate records within the allowed time frame.
- D
Enable Delta Lake change data feed and configure the sink to automatically deduplicate incoming records at the storage layer.
Why it fails: Delta Lake change data feed tracks row-level changes for downstream consumers rather than automatically deduplicating incoming streaming micro-batches, meaning duplicate records from upstream sources will still be persisted unless handled in the Spark logic.