Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Deduplication in streams

Batch & Streaming · Streaming Design

Deduplication in streams

Mediumbatch-streaming-38
deduplicationidempotencystate-ttlstreaming

Question

How do you deduplicate events in a streaming pipeline?

Solution

Deduplicating events in a streaming pipeline requires upstream producers to assign a globally unique identifier to every event at creation time, such as a transaction UUID or natural business key. The streaming engine tracks these observed identifiers inside a keyed state backend bounded by a watermark or time-to-live expiration window. Because unbounded state memory is technically impossible, pipelines also deduplicate again at the sink layer using merge or upsert operations to catch late-arriving duplicates that exceed the streaming deduplication window.

Producer identity requirements

Stream deduplication cannot function without an immutable event identifier generated by the source system.

Relying on broker metadata like Kafka partition offsets, message arrival timestamps, or network ingestion times is an anti-pattern. If a publishing service encounters a network timeout and retries producing an order payload, Kafka assigns a fresh offset and a new arrival timestamp to the second attempt. The stream processor must rely on the business payload identifier, such as order_id or payment_uuid, to recognize the duplicate.

In-stream state filtering

Within the stream processing topology, an operator maintains a key-value state store of seen IDs:

SELECT
  order_id,
  customer_id,
  order_amount,
  event_time
FROM (
  SELECT
    *,
    ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY event_time ASC) AS row_num
  FROM order_stream
)
WHERE row_num = 1;

This deduplication pattern retains the earliest occurrence and filters out subsequent matches.

Key considerations for managing stream state:

  • The operator must configure a State Time-To-Live (TTL), such as retaining keys for twenty-four hours.
  • Without a strict TTL or watermark boundary, the state store retains every historical identifier forever, causing eventual memory exhaustion and job crashes.

Boundaries and sink defense

Time-bounded in-stream deduplication has an unavoidable blind spot: it cannot detect duplicates that arrive after the state TTL has expired.

If a mobile client retries an upload forty-eight hours after an outage, the stream processor has already purged the original event ID from its state and will process the event as a new record. To protect against this edge case, enforce deduplication at the destination sink by writing to tables using SQL MERGE statements or primary-key upserts, providing layered protection against duplicate records.

PreviousNext