Microbatch Deduplication Strategies for Near Real-Time Data Pipelines
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, the handbook illustrates basic SQL-based deduplication patterns that underpin this approach.
Spark (Delta Lake) implementation:
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:
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:
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:
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:
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:
// 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:
- Choose a deterministic primary key (e.g., UUID or natural business key) to serve as the deduplication anchor.
- Combine windowed aggregation with a state store for environments with bounded lateness.
- Apply an upsert/merge to a Delta Lake table for the "source of truth" layer, providing both deduplication and historic audit trails.
- Run periodic compaction (e.g., nightly) to shrink raw logs and remove lingering duplicates.
- 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
MERGEstatements 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.mdunder 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 file contains SQL-based deduplication query examples used in homework assignments. The main README.md file references the external Microbatch Hourly Deduped Tutorial repository (EcZachly/microbatch-hourly-deduped-tutorial) for complete Spark and Flink implementations.
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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →