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_fnandevent_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 recordsexamples/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)
| File | Size | Uploaded | |
|---|---|---|---|
| streamdedup-0.1.0.tar.gz | 11.7 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| 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
|