# Microbatch Deduplication Strategies for Near Real-Time Data Pipelines

> Master microbatch deduplication strategies for near real-time data pipelines. Learn primary-key upserts, time-windowed aggregations, and state-store lookups for exactly-once processing.

- Repository: [DataExpert.io/data-engineer-handbook](https://github.com/DataExpert-io/data-engineer-handbook)
- Tags: deep-dive
- Published: 2026-08-07

---

**Microbatch deduplication strategies eliminate duplicate events in near real-time pipelines by applying primary-key upserts, time-windowed aggregations, or state-store lookups during ingestion to ensure exactly-once processing semantics.**

Near real-time (or "microbatch") pipelines ingest data in short, fixed intervals ranging from seconds to a few minutes. Because the same event can be delivered multiple times due to retries, out-of-order delivery, or source-side duplication, each microbatch must apply deduplication before downstream processing. The **DataExpert-io/data-engineer-handbook** repository outlines several architectural patterns for handling these duplicates in modern data engines.

## Why Microbatch Deduplication Matters

Duplicate events corrupt aggregates, inflate counts, and lead to misleading analytics in BI dashboards. Redundant processing also wastes compute resources, especially in cloud pay-as-you-go environments. Many downstream consumers assume exactly-once semantics; deduplication at the ingestion layer guarantees that contract.

## Six Proven Microbatch Deduplication Strategies

The handbook identifies six primary approaches, each with specific trade-offs for latency, throughput, and state management.

### Primary-Key Upsert (MERGE)

Each record contains a stable primary key (e.g., `event_id`). The microbatch writes to a table using a **MERGE** statement (or `INSERT … ON CONFLICT DO UPDATE`) that overwrites any prior row with the same key.

Use this for low-volume streams where key cardinality is manageable and you need the most recent state. In [`intermediate-bootcamp/materials/2-fact-data-modeling/homework/homework.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/2-fact-data-modeling/homework/homework.md), the handbook illustrates basic SQL-based deduplication patterns that underpin this approach.

**Spark (Delta Lake) implementation:**

```python
from delta.tables import DeltaTable
from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("MicrobatchDedup").getOrCreate()

source_df = spark.read.format("json").load("s3://bucket/microbatch")
target = DeltaTable.forPath(spark, "s3://bucket/delta/events")

target.alias("t").merge(
    source_df.alias("s"),
    "t.event_id = s.event_id"
).whenMatchedUpdateAll() \
 .whenNotMatchedInsertAll() \
 .execute()

```

The `merge()` method with `whenMatchedUpdateAll()` guarantees that only the latest version of each `event_id` survives, providing exactly-once semantics for downstream queries.

### Windowed Aggregation

Group records by a time-based window (e.g., 5 minutes) and unique identifier, then keep only the first or last occurrence. This works well when ingestion delay is bounded to a few minutes.

Use this for high-throughput streams where out-of-order arrivals are limited.

**Flink implementation:**

```java
DataStream<Event> stream = env
    .addSource(new FlinkKafkaConsumer<>("events", new EventDeserializationSchema(), props))
    .assignTimestampsAndWatermarks(WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(30))
        .withTimestampAssigner((e, ts) -> e.timestamp));

DataStream<Event> deduped = stream
    .keyBy(e -> e.eventId)
    .window(TumblingEventTimeWindows.of(Time.minutes(5)))
    .reduce((e1, e2) -> e1.timestamp > e2.timestamp ? e1 : e2);

```

By keying on `eventId` and applying a tumbling event-time window, duplicate events arriving within the 5-minute window collapse to the most recent record.

### Deduplication Table (State Store)

Maintain a lightweight key-value store (e.g., Redis, DynamoDB, or a Delta Lake "dedup" table) recording the latest processed `event_id`. Each microbatch checks the store before processing.

This strategy suits stateless microbatch jobs (e.g., serverless) that require fast lookups.

**Python with Redis:**

```python
import redis

r = redis.StrictRedis(host='redis', port=6379, db=0)

def process_batch(batch):
    for record in batch:
        event_id = record["event_id"]
        if r.setnx(event_id, 1):
            handle(record)
        else:
            continue

```

The `setnx()` (set-if-not-exists) command provides O(1) deduplication checks with minimal latency.

### Idempotent Writes (Append-Only + Compaction)

Write every record to an immutable log (e.g., Kafka topic, Delta Lake append) and run periodic compaction jobs that remove duplicates using row-level deduplication.

Use this when write-path latency must stay minimal and you can tolerate eventual consistency.

**Delta Lake pattern:**

```sql
INSERT INTO raw_events VALUES ...;

-- Later compaction job:
OPTIMIZE raw_events ZORDER BY (event_id) 
WHERE _insertion_timestamp < current_timestamp() - interval '1 hour'

```

### Probabilistic Bloom Filter

Store a Bloom filter of recently seen keys (e.g., last N minutes). Incoming events that hit the filter are dropped; rare false positives are acceptable.

This strategy fits very high-volume streams where exact deduplication is too expensive.

**Spark Structured Streaming:**

```scala
val filter = BloomFilter.create(1000000, 0.01)
df.filter(row => !filter.mightContain(row.getAs[String]("event_id")))

```

### Two-Phase Commit with Idempotent Tokens

The source embeds an idempotency token; the sink writes the token to a transaction log before committing. If the token exists, the transaction is ignored.

This works best with APIs that already provide idempotency tokens (e.g., Stripe webhooks).

**Flink stateful processing:**

```java
// Within KeyedProcessFunction
ValueState<Boolean> seenState = getRuntimeContext().getState(
    new ValueStateDescriptor<>("token-seen", Types.BOOLEAN)
);

public void processElement(TokenEvent event, Context ctx, Collector<Output> out) {
    if (seenState.value() == null) {
        seenState.update(true);
        out.collect(process(event));
    }
    // If state exists, token is duplicate; skip
}

```

## Architectural Recommendations

According to the DataExpert-io/data-engineer-handbook source code, implement these best practices when designing your pipeline:

1. **Choose a deterministic primary key** (e.g., UUID or natural business key) to serve as the deduplication anchor.
2. **Combine windowed aggregation with a state store** for environments with bounded lateness.
3. **Apply an upsert/merge to a Delta Lake table** for the "source of truth" layer, providing both deduplication and historic audit trails.
4. **Run periodic compaction** (e.g., nightly) to shrink raw logs and remove lingering duplicates.
5. **Instrument metrics** (duplicate-rate, latency) to monitor strategy effectiveness.

The handbook references a dedicated tutorial on microbatch deduplication: *Microbatch Hourly Deduped Tutorial* available at `EcZachly/microbatch-hourly-deduped-tutorial`. This repository contains concrete Spark and Flink examples showcasing the MERGE-based approach with state-store fallback.

## Summary

- **Microbatch deduplication** is essential for exactly-once semantics in near real-time pipelines processing data in seconds-to-minutes intervals.
- **Primary-key upserts** via `MERGE` statements provide the strongest consistency guarantees for moderate volume streams.
- **Windowed aggregations** handle high-throughput scenarios with bounded out-of-order delays using event-time processing.
- **State stores** like Redis enable stateless microbatch functions to perform O(1) duplicate checks.
- **Append-only patterns** with periodic compaction minimize write latency while ensuring eventual consistency.
- The **DataExpert-io/data-engineer-handbook** documents these patterns in [`README.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/README.md) under the Design Patterns section, with additional SQL examples in the intermediate bootcamp materials.

## Frequently Asked Questions

### What is microbatch deduplication in data engineering?

Microbatch deduplication is the process of identifying and removing duplicate events from data streams that are processed in small, discrete batches (typically seconds to minutes apart). This prevents repeated records—caused by network retries, at-least-once delivery guarantees, or producer-side errors—from corrupting analytics aggregates or triggering duplicate business actions.

### How does the MERGE strategy handle duplicates in Delta Lake?

The MERGE strategy uses ACID transactions to atomically insert new records while updating existing ones based on a primary key. In `MERGE INTO` statements, when a matching `event_id` is found, the operation updates the existing row; otherwise, it inserts the new row. This upsert pattern, documented in the handbook's design patterns, ensures idempotent writes even if the same microbatch is reprocessed due to job failures.

### When should I use a Bloom filter versus a state store for deduplication?

Use a **Bloom filter** when processing extremely high-volume streams where approximate deduplication is acceptable and memory constraints prevent storing complete key histories; the probabilistic structure provides constant-space efficiency with a configurable false-positive rate. Use a **state store** (like Redis or DynamoDB) when exact deduplication is required and you need to maintain arbitrary retention windows, accepting the operational overhead of external storage lookups and potential network latency.

### What file in the DataExpert-io/data-engineer-handbook contains deduplication examples?

The [`intermediate-bootcamp/materials/2-fact-data-modeling/homework/homework.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/2-fact-data-modeling/homework/homework.md) file contains SQL-based deduplication query examples used in homework assignments. The main [`README.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/README.md) file references the external *Microbatch Hourly Deduped Tutorial* repository (`EcZachly/microbatch-hourly-deduped-tutorial`) for complete Spark and Flink implementations.