A data pipeline can fail after writing a record but before acknowledging the message that carried it. The producer retries. The same event arrives again.
If the consumer inserts a new row every time, a transient failure becomes duplicate data.
This problem is especially important in financial-data systems, where repeated ingestion can distort totals, create confusing audit trails, or trigger downstream work more than once. The solution is not to assume that messages arrive exactly once. It is to design consumers so that processing the same logical event repeatedly has the same effect as processing it once.
This tutorial demonstrates a small, testable pattern using Python and SQL. The records are synthetic and the example is not tied to any particular financial provider.
What idempotency means
An operation is idempotent when repeating it does not change the final result after the first successful application.
For ingestion, the important question is: how does the system recognise that two deliveries represent the same logical event?
A timestamp alone is usually a poor key. Two valid events can share a timestamp, and the same event may be delivered with a different ingestion timestamp.
Prefer a stable source event identifier when the producer guarantees its uniqueness. If none exists, the consumer may need a carefully defined composite key based on stable source fields. Hashing an arbitrary payload can help detect exact duplicates, but it does not automatically identify semantically equivalent records.
1. Define the event contract
Start with a small event model:
from dataclasses import dataclass
from decimal import Decimal
@dataclass(frozen=True)
class AccountBalanceEvent:
source: str
source_event_id: str
account_key: str
currency: str
balance: Decimal
source_event_time: str
The account_key should be an internal identifier or a controlled reference—not a raw account number in logs. The source event ID should be stable across retries.
In production, validate the event before attempting a write. Check required fields, supported currency codes, numeric precision, timestamp formats, and any source-specific constraints. Reject malformed records into a controlled error path instead of silently modifying them.
2. Put uniqueness in the database
Application-level checks such as “look up the event, then insert if it does not exist” are vulnerable to race conditions. Two workers can both perform the lookup before either inserts.
A database uniqueness constraint is the durable guardrail:
CREATE TABLE ingested_balance_events (
source TEXT NOT NULL,
source_event_id TEXT NOT NULL,
account_key TEXT NOT NULL,
currency CHAR(3) NOT NULL,
balance NUMERIC(18, 2) NOT NULL,
source_event_time TIMESTAMPTZ NOT NULL,
received_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (source, source_event_id)
);
The pair (source, source_event_id) is the idempotency key in this example. Confirm that the producer's event ID really is unique within that source. If the contract differs, change the key rather than assuming uniqueness.
The database—not just the Python process—now enforces the rule across concurrent workers.
3. Use an atomic insert
For PostgreSQL, ON CONFLICT DO NOTHING provides a concise way to ignore a repeated key:
INSERT INTO ingested_balance_events (
source,
source_event_id,
account_key,
currency,
balance,
source_event_time
)
VALUES (
%(source)s,
%(source_event_id)s,
%(account_key)s,
%(currency)s,
%(balance)s,
%(source_event_time)s
)
ON CONFLICT (source, source_event_id) DO NOTHING;
If the event is new, a row is inserted. If the same event is delivered again, the existing row remains unchanged.
However, DO NOTHING is not a complete data-quality strategy. If a producer reuses an event ID for a different payload, the conflict can hide a source-contract violation. For high-integrity pipelines, compare the existing record with the incoming event and flag a mismatched payload for investigation.
Do not silently overwrite the original event just because the key collides.
4. Distinguish duplicates from conflicting replays
A duplicate delivery and an inconsistent replay are different cases:
- Exact replay: same event key and same canonical content.
- Conflicting replay: same event key but different content.
- New event: a previously unseen event key.
A useful implementation can store a canonical payload hash or compare selected normalized fields. Canonicalization must be deterministic: key order, decimal formatting, timestamps and optional fields should be handled consistently.
For example, a hash can be used as a compact comparison aid:
import hashlib
import json
def canonical_hash(payload: dict) -> str:
encoded = json.dumps(
payload,
sort_keys=True,
separators=(",", ":"),
ensure_ascii=False,
).encode("utf-8")
return hashlib.sha256(encoded).hexdigest()
This simple helper assumes that values have already been normalized into JSON-compatible forms. For financial values, convert Decimal to a defined string representation before hashing. Do not rely on a hash of unnormalized input if "10.0" and "10.00" are intended to mean the same amount.
A hash is not a substitute for validation, and it should not be treated as encryption or as a way to anonymize sensitive data.
5. Keep ingestion and downstream effects consistent
A common failure pattern is:
- Insert the event.
- Publish a notification.
- Crash before acknowledging the message.
- Receive the message again and publish a second notification.
Deduplicating the database row does not automatically deduplicate every side effect.
One common approach is the transactional outbox pattern. Within the same database transaction, insert the accepted event and an outbox record describing the downstream work. A separate publisher sends outbox records and marks them as delivered. The publisher itself must tolerate retries, because a crash can still occur after sending but before recording success.
A simplified outbox table might look like this:
CREATE TABLE outbox (
outbox_id BIGSERIAL PRIMARY KEY,
event_key TEXT NOT NULL UNIQUE,
event_type TEXT NOT NULL,
payload JSONB NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
published_at TIMESTAMPTZ
);
Use a stable event_key to prevent creating the same logical outbox event more than once. Consumers downstream should also be idempotent where possible.
This pattern does not magically guarantee end-to-end exactly-once delivery. It makes the system's retry behaviour explicit and helps achieve effectively-once outcomes for defined operations.
6. Test the failure paths
The happy path is not enough. Test at least these cases:
- The first delivery inserts one row.
- Replaying the same event leaves the row count unchanged.
- Concurrent deliveries cannot create duplicate keys.
- Reusing an event ID with different content is detected.
- Invalid amounts and timestamps are rejected.
- A database failure can be retried safely.
- An outbox publisher can restart without creating duplicate downstream effects.
A conceptual test for the core property is:
def test_replaying_event_does_not_duplicate_it(repository, event):
repository.ingest(event)
repository.ingest(event)
rows = repository.find_by_event_key(
event.source,
event.source_event_id,
)
assert len(rows) == 1
The repository and fixture are intentionally abstract; adapt the test to your database layer and test environment. Use a real database integration test for uniqueness and concurrency behaviour, because an in-memory mock cannot verify database constraints.
7. Monitor more than successful inserts
A pipeline can report success while silently dropping conflicting replays. Track metrics such as:
- New events accepted
- Exact duplicates ignored
- Conflicting event-key replays
- Validation failures
- Database retry counts
- Outbox backlog and delivery latency
- Age of the oldest unprocessed event
Avoid placing raw personal or account data in metric labels, logs or tracing attributes. Use controlled identifiers and aggregate metrics wherever possible.
A compact design checklist
Before shipping an ingestion consumer, confirm that:
- The idempotency key is stable and documented.
- The database enforces uniqueness atomically.
- Duplicate and conflicting replays are handled differently.
- Validation happens before persistence.
- Downstream side effects have their own retry strategy.
- Tests cover concurrency and failure recovery.
- Logs and metrics avoid unnecessary sensitive data.
- The event contract and key strategy are versioned when changed.
Conclusion
Retries are normal in distributed systems. Duplicate effects do not have to be.
A stable event identity, database-enforced uniqueness, explicit conflict handling, and idempotent downstream processing create a much safer foundation than relying on the consumer to receive every message exactly once. Start with the event contract and database constraint, then test the failure modes that can occur between each step.
AI assistance disclosure: This draft was prepared with AI assistance. Review the design, validate the code against your stack, and add your own experience before publishing.
Top comments (0)