Data quality and data contracts
Keeping data pipelines trustworthy: common failure modes, data contracts between producers and consumers, schema and semantic validation, freshness, volume and distribution checks, quarantining bad records, lineage, backfills, and alerting on data like you alert on services.
Reading is half of it. See this used in a real interview: walk through Design an Ad Click Aggregator →
Services fail loudly; data pipelines fail silently. A producer renames a field, a mobile release stops sending an event, a currency column switches from cents to dollars, and dashboards, billing and machine learning models quietly go wrong for days. Data-heavy designs (ad aggregation, analytics, recommendations) are judged partly on how they prevent and detect this. This guide covers the practices that keep data trustworthy.
How data goes wrong
- Schema changes: renamed, removed or retyped fields break consumers.
- Semantic changes: same field, different meaning (units, time zones, definitions of "active user").
- Missing data: a source stops sending, a job silently processes zero rows, a partition is empty.
- Duplicates: retries and replays double-count. See delivery semantics.
- Late data: events arrive after aggregates were computed. See windowing and watermarks.
- Bad values: nulls, negative prices, impossible timestamps, test traffic in production data.
- Logic bugs in transformations.
Data contracts
A data contract is an explicit agreement between a producer and its consumers about a dataset or event stream:
- Schema: fields, types, required versus optional.
- Semantics: what each field means, units, allowed values.
- Quality guarantees: completeness, freshness (available within N minutes), uniqueness keys.
- Ownership: who to contact, and the change process (notice periods, versioning).
Enforce the schema part automatically: producers validate against the contract before publishing, and a schema registry rejects incompatible changes. See schema evolution and serialization. Treat events and tables used by other teams as public APIs, not internal implementation details.
Checks in the pipeline
Validate data at each stage, like tests for data:
| Check | Example |
|---|---|
| Schema | required columns present, types correct |
| Nulls and ranges | price >= 0, country in the allowed list, no null ids |
| Uniqueness | one row per (order_id) |
| Referential | every ad_id exists in the ads table |
| Freshness | latest partition no older than 1 hour |
| Volume | today's row count within 30 % of the same weekday last week |
| Distribution | share of mobile traffic not suddenly 0 % or 100 % |
| Reconciliation | sum of revenue matches the payments system within tolerance |
Run cheap checks on every batch or micro-batch, and heavier ones daily.
When checks fail
- Block: do not publish the output partition; downstream consumers keep yesterday's data rather than wrong data. Use write-audit-publish: write to a staging location, run checks, then atomically publish.
- Quarantine: route bad records to a separate table for inspection, process the rest.
- Alert: the owning team gets paged for critical datasets, like a service outage.
Which to choose depends on the cost of wrong data versus late data. Billing data should block; an internal dashboard may quarantine and continue.
Lineage
Lineage records which datasets and jobs produce which outputs. When a source breaks, lineage shows every downstream table, dashboard and model affected, so you can notify owners and plan backfills. Catalogs and orchestrators capture it automatically. See data lakes and lakehouses.
Backfills and reprocessing
After fixing a bug or a bad source, recompute affected outputs:
- Keep raw data immutable and long enough to reprocess.
- Make jobs idempotent per partition (overwrite a day's output), so rerunning is safe.
- Use the orchestrator to backfill date ranges in order of dependency. See workflow orchestration.
- Communicate corrections to consumers (especially for billing and reported metrics).
Observability for data
Track datasets like services: freshness, volume, failed checks and job duration, with SLOs for critical ones ("billing table complete by 06:00 UTC, 99.5 % of days"). See SLIs, SLOs and error budgets. Anomaly detection on volumes and distributions catches problems no one wrote a rule for.
Machine learning data
Models amplify data problems: a broken feature silently degrades predictions. Validate training data and serving features, monitor feature distributions for drift, and keep training and serving feature definitions shared. See feature stores and ML serving.
In the interview
For Design an Ad Click Aggregator: "Click events follow a contract enforced by the schema registry; the stream job deduplicates by event id; each hourly aggregate is written to staging, checked for volume against last week and reconciled against raw counts, then published; failures block billing outputs and page the owning team; raw events are retained for backfills." For Design a Monitoring System, mention alerting when a metric source goes silent.
Checklist
- Contracts with schema, semantics, guarantees and owners.
- Schema enforcement at the producer and registry.
- Checks for nulls, ranges, uniqueness, freshness, volume, distribution and reconciliation.
- Write-audit-publish, quarantine or block depending on the cost of errors.
- Lineage for impact analysis.
- Immutable raw data and idempotent, partitioned jobs for backfills.
- Data SLOs and anomaly alerts.