# How to Implement Data Deduplication Strategies for Microbatch Processing

> Learn how to implement data deduplication strategies for microbatch processing using stateful filtering, window functions, and idempotent writes for exactly-once semantics.

- Repository: [DataExpert.io/data-engineer-handbook](https://github.com/DataExpert-io/data-engineer-handbook)
- Tags: how-to-guide
- Published: 2026-08-12

---

**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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/funnel_analysis.sql), the handbook demonstrates a common table expression (CTE) pattern that assigns row numbers to events partitioned by unique identifiers:

```sql
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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-applying-analytical-patterns/lecture-lab/funnel_analysis.sql) ([line 1](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-applying-analytical-patterns/lecture-lab/funnel_analysis.sql#L1)), 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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/player_game_edges.sql) ([line 2](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/1-dimensional-data-modeling/lecture-lab/player_game_edges.sql#L2)), 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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/team_vertex_job.py) ([line 5](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/3-spark-fundamentals/src/jobs/team_vertex_job.py#L5)) demonstrates this pattern in PySpark:

```python
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:

```sql
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.

## Advanced State Management with Flink

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

```java
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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/funnel_analysis.sql) ([line 1](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-applying-analytical-patterns/lecture-lab/funnel_analysis.sql#L1)).
- **Stateful deduplication** leverages watermarks and checkpointed state stores to remember keys across microbatches indefinitely, implemented in [`team_vertex_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/team_vertex_job.py) ([line 5](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/3-spark-fundamentals/src/jobs/team_vertex_job.py#L5)).
- **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](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/README.md#L340)) 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.