# How to Implement Microbatch Deduplication Patterns in Data Engineering Pipelines

> Learn microbatch deduplication patterns in data engineering pipelines using SQL, Spark, or Delta Lake. Achieve exactly-once semantics for streaming data.

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

---

**Microbatch deduplication eliminates duplicate records within each micro-batch using SQL GROUP BY aggregations, Spark SQL window functions, or Delta Lake MERGE operations, ensuring exactly-once semantics for high-frequency streaming pipelines.**

Microbatch deduplication is a fundamental pattern for building reliable, idempotent data pipelines that ingest event-level data at high frequency. By removing duplicates before downstream transformations run, you guarantee that aggregates, slowly changing dimensions, and windowed analytics operate on a clean, unique event set. This guide examines production-ready implementations from the **DataExpert-io/data-engineer-handbook** repository, demonstrating how to encode these patterns in SQL, Spark, and Delta Lake.

## Why Microbatch Deduplication Matters

Without deduplication, streaming pipelines suffer from over-counted metrics, corrupted state, and exploded join cardinalities. When Kafka or Kinesis delivers the same event multiple times, or when you re-process a microbatch after a failure, duplicate rows propagate through your pipeline and distort business logic.

Implementing deduplication at the microbatch level provides three critical benefits:

- **Deterministic results**: Each processing window contains exactly one occurrence of every unique event, preventing double-counting in funnel and retention calculations.
- **Safe replay**: Re-reading the same microbatch after a failure becomes a no-op, eliminating manual cleanup tasks.
- **Optimized joins**: Removing duplicate keys before joins prevents cardinality explosions and reduces shuffle overhead in distributed processing.

## Core Architectural Strategies for Microbatch Deduplication

The Data Engineer Handbook demonstrates five primary approaches to microbatch deduplication, each suited to specific tooling and latency requirements.

**SQL GROUP BY deduplication** aggregates on natural business keys (`event_id`, `user_id`, `timestamp`) to retain the first or latest occurrence. This works universally across Snowflake, Redshift, BigQuery, and Postgres.

**Window functions (ROW_NUMBER)** assign a sequential index per key ordered by event time, filtering to `row_number = 1` to keep only the latest record. This pattern executes efficiently in Spark SQL and cloud data warehouses.

**Idempotent upserts** use Delta Lake or BigQuery MERGE statements with primary keys, converting duplicate inserts into no-ops. This approach maintains state across microbatches and handles slowly changing dimensions.

**Watermark-based state cleanup** in Flink or Structured Streaming automatically drops old duplicates by expiring state outside the watermark threshold.

**Hash-based deduplication** computes deterministic hashes of payloads in Kafka Streams or Flink, dropping records with hashes already observed in the current microbatch.

## Repository Implementation Examples

The following examples from the DataExpert-io/data-engineer-handbook repository illustrate how to implement these patterns in production code.

### SQL GROUP BY Deduplication for Funnel Analysis

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), the pipeline uses a common table expression to deduplicate events before calculating conversion rates. The `deduped_events` CTE groups on all event attributes to retain one row per unique combination of `url`, `host`, `user_id`, and `event_time`.

```sql
WITH deduped_events AS (
    SELECT
        url,
        host,
        user_id,
        event_time
    FROM events
    GROUP BY 1, 2, 3, 4
),
clean_events AS (
    SELECT *, DATE(event_time) AS event_date
    FROM deduped_events
    WHERE user_id IS NOT NULL
    ORDER BY user_id, event_time
),
converted AS (
    SELECT
        ce1.user_id,
        ce1.event_time,
        ce1.url,
        COUNT(DISTINCT CASE WHEN ce2.url = '/api/v1/user' THEN ce2.url END) AS converted
    FROM clean_events ce1
    JOIN clean_events ce2
      ON ce2.user_id = ce1.user_id
     AND ce2.event_date = ce1.event_date
     AND ce2.event_time > ce1.event_time
    GROUP BY 1, 2, 3
)
SELECT url,
       COUNT(1) AS impressions,
       CAST(SUM(converted) AS REAL) / COUNT(1) AS conversion_rate
FROM converted
GROUP BY 1
HAVING CAST(SUM(converted) AS REAL) / COUNT(1) > 0
   AND COUNT(1) > 100;

```

This approach ensures that all downstream funnel calculations operate on a clean event stream without duplicate page views inflating impression counts.

### Spark SQL ROW_NUMBER Deduplication

For scenarios requiring the latest record per entity, [`intermediate-bootcamp/materials/1-dimensional-data-modeling/lecture-lab/team_vertices.sql`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/1-dimensional-data-modeling/lecture-lab/team_vertices.sql) demonstrates the window function pattern. The code partitions by the natural key (`team_id`) and orders by `event_time` descending, retaining only the most recent row.

```scala
val teams = spark.read.table("raw_teams")
val deduped = teams
  .withColumn("rn",
    row_number()
      .over(Window.partitionBy($"team_id")
                 .orderBy($"event_time".desc))
  )
  .filter($"rn" === 1)
  .drop("rn")
deduped.createOrReplaceTempView("teams_deduped")

```

The `row_number` function assigns a unique index to each duplicate within the microbatch. Filtering to `rn = 1` efficiently discards older duplicates while preserving the complete latest record.

### Idempotent Upserts with Delta Lake

The [`intermediate-bootcamp/materials/3-spark-fundamentals/src/jobs/players_scd_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/3-spark-fundamentals/src/jobs/players_scd_job.py) file implements microbatch deduplication using Delta Lake's MERGE operation. This pattern is critical for slowly changing dimension (SCD) processing, where the same `(player_id, event_time)` pair might appear in multiple microbatches during pipeline replays.

```python
from delta.tables import DeltaTable

deltaTable = DeltaTable.forPath(spark, "/mnt/delta/players")
new_events = spark.read.parquet("/mnt/raw/player_events")

deltaTable.alias("t").merge(
    new_events.alias("s"),
    "t.player_id = s.player_id AND t.event_time = s.event_time"
).whenNotMatchedInsertAll().execute()

```

The MERGE statement matches on the composite key (`player_id`, `event_time`). When a duplicate arrives, the condition `t.player_id = s.player_id` matches an existing row, causing the `whenNotMatchedInsertAll` clause to skip the insert, effectively making the operation idempotent.

## Step-by-Step Microbatch Pipeline Implementation

Implementing microbatch deduplication requires integrating these patterns into a cohesive pipeline architecture:

1. **Ingest raw events** into a staging table such as `events_staging` with minimal transformation to maintain low latency.
2. **Apply deduplication** using the appropriate pattern for your engine—SQL GROUP BY for set-based processing, ROW_NUMBER for keeping latest records, or Delta MERGE for stateful upserts.
3. **Write clean data** to the permanent table or Delta Lake location, ensuring downstream consumers receive exactly-once semantics.
4. **Trigger downstream jobs** for funnels, retention analysis, and SCD processing, confident that the input contains no duplicates.

Because deduplication occurs within each microbatch, the pipeline maintains low latency while guaranteeing logical exactly-once processing.

## Summary

- **Microbatch deduplication** removes duplicate events within processing windows to prevent over-counting and state corruption in streaming pipelines.
- **SQL GROUP BY** provides a simple, universal deduplication mechanism suitable for batch and microbatch SQL engines.
- **ROW_NUMBER window functions** efficiently retain the latest record per natural key in Spark SQL and cloud data warehouses.
- **Delta Lake MERGE** enables idempotent upserts that handle duplicates across microbatches by treating inserts as no-ops when primary keys exist.
- The **DataExpert-io/data-engineer-handbook** repository provides production-ready examples in [`funnel_analysis.sql`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/funnel_analysis.sql), [`team_vertices.sql`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/team_vertices.sql), and [`players_scd_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/players_scd_job.py).

## Frequently Asked Questions

### What is microbatch deduplication?

Microbatch deduplication is a data processing pattern that eliminates duplicate records within discrete batches of streaming data, typically processed at intervals of seconds or minutes. Unlike stream processing with per-record operations, microbatch deduplication operates on bounded sets of data, using SQL aggregations or window functions to ensure each unique event appears exactly once before downstream transformations execute.

### How does Delta Lake MERGE handle duplicate events?

Delta Lake MERGE handles duplicates by matching incoming records against existing table rows using a specified key condition. In the [`players_scd_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/players_scd_job.py) implementation, the merge condition `t.player_id = s.player_id AND t.event_time = s.event_time` checks for existing records. When a match exists, the `whenNotMatchedInsertAll` clause prevents insertion, effectively ignoring the duplicate. This approach maintains idempotency across pipeline reruns and backfills.

### When should I use ROW_NUMBER versus GROUP BY for deduplication?

Use **ROW_NUMBER** when you need to retain specific records based on ordering criteria, such as keeping the latest event per `user_id` sorted by `event_time`. Use **GROUP BY** when you need to remove exact duplicates regardless of ordering, or when calculating aggregates where any single representative row suffices. ROW_NUMBER requires window function support (Spark SQL, modern data warehouses), while GROUP BY works universally across SQL engines.

### Can microbatch deduplication handle late-arriving data?

Microbatch deduplication handles late-arriving data only within the retention window of the current microbatch. For late data arriving in subsequent batches, you must implement idempotent upserts using Delta Lake MERGE or similar mechanisms that maintain state across batches. Watermark-based strategies in Flink or Structured Streaming automatically expire state older than the watermark threshold, dropping duplicates that arrive too late for the defined window.