How to Implement Data Deduplication Strategies for Microbatch Processing

Implement data deduplication strategies for microbatch processing by combining watermark-based stateful filtering, key-based window functions, and idempotent sink writes to achieve exactly-once semantics across distributed streams.

Microbatch processing frameworks like Apache Spark Structured Streaming and Apache Flink process data in discrete intervals, making them susceptible to duplicate records from retries, out-of-order arrivals, or upstream system failures. The DataExpert-io/data-engineer-handbook provides practical implementations of data deduplication strategies for microbatch processing, demonstrating how to maintain data integrity without sacrificing throughput. These patterns range from simple SQL window functions to stateful stream processing with TTL-managed key stores.

Key-Based Window Deduplication with SQL

The simplest approach to deduplication uses deterministic keys within time-bounded windows. In funnel_analysis.sql, the handbook demonstrates a common table expression (CTE) pattern that assigns row numbers to events partitioned by unique identifiers:

WITH deduped_events AS (
    SELECT *,
           ROW_NUMBER() OVER (
               PARTITION BY event_id 
               ORDER BY event_time ASC
           ) as rn
    FROM raw_events
    WHERE event_time >= CURRENT_TIMESTAMP - INTERVAL '1' HOUR
)
SELECT * FROM deduped_events WHERE rn = 1

This pattern appears in intermediate-bootcamp/materials/4-applying-analytical-patterns/lecture-lab/funnel_analysis.sql (line 1), where ROW_NUMBER() eliminates duplicates by keeping only the first occurrence of each event_id within the processing window. A similar deduplication CTE appears in player_game_edges.sql (line 2), illustrating how edge creation pipelines handle duplicate player-game relationships.

Stateful Deduplication with Watermarks

For high-volume streams requiring cross-window memory, use state-store based deduplication. This strategy persists seen keys in a fault-tolerant state store—such as Spark’s checkpoint directory or Flink’s RocksDB backend—and filters incoming records against this historical state.

The handbook’s team_vertex_job.py (line 5) demonstrates this pattern in PySpark:

from pyspark.sql import functions as F

# Create deduplicated view using watermark and dropDuplicates

deduped_df = (input_df
    .withWatermark("event_time", "10 minutes")
    .dropDuplicates(["event_id"]))

deduped_df.createOrReplaceTempView("teams_deduped")

The withWatermark operator tells the engine to maintain state for 10 minutes beyond the maximum observed event time, allowing late arrivals to be checked against the deduplication key store before being dropped.

Idempotent Writes and Upserts

When downstream sinks support transactional semantics, implement idempotent writes as a safety net. Rather than filtering upstream, write every record and let the storage layer resolve conflicts based on primary keys.

For Delta Lake sinks, use the MERGE INTO syntax:

MERGE INTO prod.events AS target
USING (SELECT * FROM staged.events) AS source
ON target.event_id = source.event_id
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *

This complements stateful deduplication by handling ultra-late arrivals that exceed watermark boundaries. Cloud data warehouses like BigQuery offer similar INSERT ... ON CONFLICT semantics for equivalent protection.

For pure streaming contexts requiring custom logic, Flink’s KeyedProcessFunction provides fine-grained control over state TTL and deduplication logic:

public class DeduplicationFunction extends KeyedProcessFunction<String, Event, Event> {
    private ValueState<Long> lastSeen;
    
    @Override
    public void open(Configuration cfg) {
        ValueStateDescriptor<Long> desc = 
            new ValueStateDescriptor<>("lastSeen", Types.LONG);
        desc.enableTimeToLive(Time.days(1));
        lastSeen = getRuntimeContext().getState(desc);
    }
    
    @Override
    public void processElement(Event value, Context ctx, Collector<Event> out) 
            throws Exception {
        Long prev = lastSeen.value();
        if (prev == null || value.getTimestamp() > prev) {
            out.collect(value);
            lastSeen.update(value.getTimestamp());
        }
    }
}

The state TTL (time-to-live) configuration prevents unbounded growth by automatically expiring keys after 24 hours, balancing deduplication accuracy against memory constraints.

Managing Watermarks and Late Arrivals

Effective deduplication requires careful watermark configuration:

  • Watermark delay: Set withWatermark("event_time", "X minutes") to match your expected late-arrival latency. The handbook examples typically use 5–10 minute watermarks.
  • Deterministic keys: Use globally unique identifiers (UUIDs or content hashes) rather than timestamps alone, which can collide across batches.
  • Side outputs: Route records arriving after the watermark to a "late arrivals" stream for manual inspection rather than silent dropping.

Summary

  • Key-based window deduplication uses ROW_NUMBER() or dropDuplicates() within time-bounded SQL queries or DataFrame operations, as shown in funnel_analysis.sql (line 1).
  • Stateful deduplication leverages watermarks and checkpointed state stores to remember keys across microbatches indefinitely, implemented in team_vertex_job.py (line 5).
  • Idempotent sinks provide a final defense against duplicates using merge/upsert operations in Delta Lake or similar transactional storage systems.
  • TTL management prevents state store bloat by expiring old keys after a configurable duration.

Frequently Asked Questions

How does microbatch deduplication differ from batch deduplication?

Microbatch deduplication maintains state across executions to remember keys from previous batches, whereas batch deduplication processes each dataset in isolation. Batch jobs typically use simple DISTINCT or GROUP BY operations on the current dataset, while microbatch pipelines require watermarks and checkpointed state stores to handle late-arriving duplicates that span multiple processing intervals.

What watermark duration should I choose for deduplication windows?

Choose a watermark duration that balances latency against completeness. The handbook examples suggest 5–10 minutes for most use cases, though you should analyze your upstream system's retry behavior and network latency. Set the watermark to the 99th percentile of observed late-arrival times; records arriving after this threshold will be dropped or routed to side outputs.

Can I deduplicate using multiple keys simultaneously?

Yes, pass an array of columns to dropDuplicates() in Spark or use composite keys in Flink’s KeyBy selector. For SQL-based deduplication, include multiple columns in the PARTITION BY clause: PARTITION BY user_id, session_id, event_type. Ensure these combinations form a true unique identifier to avoid accidentally collapsing distinct events.

Where can I find the external tutorial referenced in the handbook?

The README.md (line 340) references an external "Microbatch Deduplication" tutorial repository (EcZachly/microbatch-hourly-deduped-tutorial) that provides additional context on hourly deduplication patterns. This resource complements the handbook’s SQL and PySpark examples with alternative implementations for specific streaming platforms.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →