Metadata-Version: 2.4
Name: streamdedup
Version: 0.1.0
Summary: Reusable streaming deduplication and late-arrival watermarking for data pipelines
Author-email: "Taiwo Hassan (Tycoach)" <davidtaiwo56@gmail.com>
License-Expression: MIT
Project-URL: Repository, https://github.com/tycoach/dedup-watermark
Keywords: data-engineering,streaming,deduplication,watermark,bloom-filter
Requires-Python: >=3.9
Description-Content-Type: text/markdown
License-File: LICENSE
Provides-Extra: redis
Requires-Dist: redis>=5.0; extra == "redis"
Provides-Extra: dev
Requires-Dist: pytest>=7.4; extra == "dev"
Requires-Dist: pytest-cov>=4.1; extra == "dev"
Dynamic: license-file

# 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:

```python
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)

```bash
pip install -e ".[dev]"
```

## Run tests

```bash
pytest tests/ -v
```

## Run the benchmark

```bash
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.
