Replayable Batch Ingestion
Separate delivery identity from business revision and make repeated input converge on a defined accepted state.
On this page
A pipeline can receive the same batch twice. It can also stop after committing data but before acknowledging delivery. The next attempt may have no trustworthy evidence that the previous process finished. A replayable design makes those ordinary conditions converge on a defined state rather than producing additional business records.
Start by separating three identities: the delivery attempt, the input batch and the business record. A retry creates a new attempt while processing the same batch. A later correction can arrive in a different batch while referring to the same business record. Using one identifier for all three hides these relationships.
Give immutable input a manifest
Record the batch identity, source, content hash, expected row count and schema revision. Treat the input bytes for that identity as immutable. If a producer sends different bytes under an existing batch identifier, report a conflict instead of quietly redefining the previous delivery.
An arrival timestamp is useful evidence but a weak identity. Two retries can arrive at different times, and a correction can arrive later while describing an earlier business state. The manifest should distinguish delivery chronology from the producer's record revision.
Retain the input for the agreed replay interval under an appropriate access policy. A pipeline that promises replay while immediately discarding its only input cannot carry out that procedure. The retained manifest alone is insufficient if the bytes it identifies are no longer available.
Validate before changing the accepted state
Parse into staging and apply the batch contract. Check required keys, units and revision types. Detect conflicting rows for the same business identity and revision inside the batch. Deduplicating them by arbitrary input order turns a producer conflict into an undocumented choice.
A current-state projection can accept only a higher valid producer revision. This PostgreSQL example illustrates the update boundary using a temporary table:
BEGIN;
CREATE TEMP TABLE current_line (
source_name text NOT NULL,
line_id text NOT NULL,
revision bigint NOT NULL,
amount numeric(12, 2) NOT NULL,
PRIMARY KEY (source_name, line_id)
);
INSERT INTO current_line VALUES ('alpha', 'A-101', 1, 12.00);
INSERT INTO current_line AS target
(source_name, line_id, revision, amount)
VALUES ('alpha', 'A-101', 2, 13.00)
ON CONFLICT (source_name, line_id) DO UPDATE
SET revision = EXCLUDED.revision, amount = EXCLUDED.amount
WHERE EXCLUDED.revision > target.revision;
SELECT * FROM current_line;
ROLLBACK;The final row has revision 2 and amount 13.00. The example leaves no persistent table. It illustrates one accepted-state rule, not a complete loader or a guarantee about delivery behavior.
Handle equal revisions deliberately
The update condition ignores a repeated revision. That is safe only after checking that an equal revision carries the same accepted payload. An equal revision with a different amount is a conflict that this small SQL statement does not detect by itself.
Keep that comparison in validation or in a separately designed database procedure. A lower revision may be an ordinary delayed delivery, but a producer reset or correction policy can require a different response. The rule must come from the producer contract rather than an assumption that every larger integer is authoritative.
For details of the conflict clause, consult PostgreSQL's INSERT documentation. Application-level identity and conflict decisions remain necessary around the database statement.
Commit data and completion together
Where they share one transactional database, commit the accepted changes and the batch completion record in the same transaction. A retry can then inspect that completion record and verify the manifest before deciding whether work is already complete.
External side effects need their own boundary. Sending an email or publishing an object is not automatically reversed by rolling back the database transaction. A transactional outbox or an idempotent destination can be part of a larger design, but each needs explicit identity and retry semantics.
Avoid claiming exactly-once execution simply because one insert has a unique key. The process may execute repeatedly. The useful claim is narrower: repeated accepted input converges on the intended stored result under the stated keys, revision rules and transaction boundaries.
Test interruption at meaningful points
Use a disposable database and an immutable fixture batch. Interrupt before the transaction commits, then replay. The accepted state should match a clean run. Next interrupt after commit but before delivery acknowledgement. Replaying should recognize the completed batch rather than create another business result.
Replay the same batch twice without interruption, deliver an older revision after a newer one, and deliver a changed payload under the same batch identity. These cases distinguish retry handling from conflict handling. A test containing only new records cannot establish either property.
Record the final accepted keys, revisions and completion state. Keep attempt diagnostics separate from business totals so that retries do not inflate accounting. Replayability is achieved when input identity, validation, accepted-state rules and completion evidence agree on what a repeated delivery means.