Skip to main content

streamdedup

Reusable streaming deduplication and late-arrival watermarking for data pipelines.

Given a stream of records, streamdedup decides — for each one — whether it's a new record, a duplicate, or arriving too late to trust:

from datetime import timedelta
from streamdedup import DedupWatermarkProcessor

processor = DedupWatermarkProcessor(
    key_fn=lambda r: r["trip_id"],
    event_time_fn=lambda r: r["pickup_datetime"],
    watermark_delay=timedelta(hours=2),
    dedup_backend="bloom",   # or "exact"
)

decision = processor.process(record)  # ACCEPT / DUPLICATE / LATE_DROPPED / LATE_ACCEPTED

Why this exists

Every pipeline that ingests events or records eventually has to answer the same two questions: have I seen this before? and is this arriving too late to still count? Retries, at-least-once delivery, multiple producers, and network delays make both questions unavoidable at any real scale.

This library answers them once, generically, so it can be dropped into any pipeline rather than reimplemented per project.

Design

  • Schema-agnostic: the processor never touches your record's shape directly — you supply key_fn and event_time_fn, two small functions. No pipeline-specific code lives inside the library.
  • Two dedup backends behind one interface (DedupBackend):
    • exact — hash-set based, zero false positives, memory scales with distinct-key cardinality in the retention window.
    • bloom — time-bucketed Bloom filter, fixed memory footprint regardless of volume, tunable false-positive rate, no false negatives.
  • Bounded memory: both backends evict state for time buckets that have fallen behind the watermark, so long-running streams don't grow unbounded.
  • No I/O inside the library: it doesn't know about Kafka, files, or databases — feed it records from anywhere, get decisions back.

Install (local development)

pip install -e ".[dev]"

Run tests

pytest tests/ -v

Run the benchmark

python benchmarks/bench_exact_vs_bloom.py

Examples

  • examples/example_tlc_usage.py — wiring against NYC TLC trip records
  • examples/example_generic_usage.py — the same library against a completely different (insurance-claims-shaped) schema, to demonstrate portability

Upgrade path

Dedup state currently lives in process memory. For a pipeline that needs dedup state to survive a restart (e.g. a production claims pipeline), a Redis-backed DedupBackend implementation can be dropped in without changing DedupWatermarkProcessor at all — see state/redis_store.py for the sketch.

Metadata

Release files for streamdedup 0.1.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for streamdedup 0.1.0
File Size Uploaded
streamdedup-0.1.0.tar.gz 11.7 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for streamdedup 0.1.0
File Interpreter ABI Platform
streamdedup-0.1.0-py3-none-any.whl Python 3 none any Details

Total release size: 22.5 kB

Release files / streamdedup-0.1.0.tar.gz

Download URL streamdedup-0.1.0.tar.gz
Size 11.7 kB
Tags Source
SHA-256 checksum
How to use checksums
ab8fe63285984b293befb90eca483c7a544e8a5c751d5d2802f8d18c05d4cfa1
BLAKE2b-256 checksum
How to use checksums
859b5cc0d2c0b5fec1503b06bd208350af2a1ea7179ec741dc55adef77a5e5cd
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.14.5

Release files / streamdedup-0.1.0-py3-none-any.whl

Download URL streamdedup-0.1.0-py3-none-any.whl
Size 10.8 kB
Tags Python 3
SHA-256 checksum
How to use checksums
656e07f3d386198c2e2473b487d5d8cdaf32daa51b76fc6468de79331a4b0920
BLAKE2b-256 checksum
How to use checksums
8fd4712329b675ecf614eba2b1701f38df35f23e60851efb45b5bc77a786aee1
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.14.5

Release history Release notifications | RSS feed

This release

0.1.0 This release

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page