A source contract fixes record identity, event time, schema expectations and ownership before a pipeline edits or aggregates incoming data.
Data source contracts: preserve raw records before transformation
Name the source boundary
A receipt processor emits submissions and later corrections. The ingest job must know which field is the stable receipt ID, which timestamp describes the business event, which timestamp describes arrival, and who can change the schema. Capture raw bytes or a governed immutable copy before casting fields. Without the raw record, a transformation bug cannot be replayed after the source has moved on. Analysis snapshots] depend on this earlier boundary.
Validate but do not erase
Reject a record with a missing ID or an impossible timestamp into a quarantine stream with its source offset and reason. Do not quietly drop it, and do not coerce a malformed amount to zero. A contract may permit additive optional fields while rejecting a changed key type. Version the rule and count each rejection reason. An alert should compare the rejection rate with the source volume; one bad row in 47 records is very different from one in millions.
Retain provenance
Store source system, partition or file identity, record position, ingest time and contract version. If an upstream team republishes a file, identify whether it is the same bytes or a new revision. Separate raw storage permissions from curated access; the raw layer may contain sensitive fields analysts should not see. Quality reconciliation] uses the provenance to trace a missing receipt.
Prove a schema change
Run the ingest contract against an old valid record, an additive-field record, a missing-key record and a changed timestamp format. The first two should be accepted under the written rule; the latter two should be quarantined with precise reasons. A parse success is not proof of semantic validity.
Implementation
from dataclasses import dataclass
from datetime import datetime
@dataclass(frozen=True)
class ReceiptEvent:
receipt_id: str
event_time: datetime
amount_cents: int
def validate_receipt(record: ReceiptEvent) -> None:
if not record.receipt_id or record.event_time.tzinfo is None:
raise ValueError("receipt identity and timezone are required")
if record.amount_cents < 0:
raise ValueError("negative receipt amount")Performance and operating cost
Validation is O(N) for N records and usually small per row. Retaining raw input costs storage and access-control work; it is also what makes a correct replay possible.
Common Mistakes
- Do not discard raw records after a failed transform.
- Do not use arrival time as event time without saying so.
- Do not coerce malformed money to a valid-looking zero.
Read next
- Data quality gates: quarantine bad rows and reconcile complete batches
- Incremental extraction: advance a compound watermark without losing ties
- Reproducible analysis snapshots: pin data, code and cutoff together
- Missing data policy: distinguish absence from a measured zero
Continue the workflow: Schema compatibility and consumer rollout.
Continue the workflow: Data classification and access boundaries.
Continue the workflow: Contract-to-test compilation and negative fixtures.
