how to detect schema drift in a data pipeline
Shows how to detect schema drift by snapshotting table schemas and diffing every run. Use when upstream producers change columns without warning, when pipelines break on unexpected schema changes, or when setting up data contracts. Not for data-value drift (different problem), for one-off schema migrations you control, or for unstructured data.
TL;DR
Snapshot the schema (column names and types) on every run and diff it against the previous snapshot; alert on any change. The fix is a small schema-check task that queries information_schema, stores the result, and fails or warns when columns are added, dropped, renamed, or retyped. Upstream producers change schemas without telling you; this is how you find out before your pipeline does.
how to detect schema drift in a data pipelineUse this when
- An upstream table changed columns and broke your load
- You want early warning before schema changes hit production
- You are setting up data contracts between producers and consumers
Not for this skill when
- The data values drifted but the schema is unchanged (data drift, different skill)
- You control the migration and just need to apply it (use migrations)
- The source is unstructured (JSON blobs, free text)
Steps
- Capture a baseline schema snapshot into a table or file:
CREATE TABLE schema_snapshots AS
SELECT table_name, column_name, data_type, ordinal_position,
CURRENT_DATE AS snapshot_date
FROM information_schema.columns
WHERE table_schema = 'staging';Expected output: a baseline recording every column and type. This is the "before" picture everything gets compared against.
- Add a task that re-captures the schema each run and diffs it against the latest snapshot:
WITH current AS (
SELECT table_name, column_name, data_type FROM information_schema.columns
WHERE table_schema = 'staging'
),
previous AS (
SELECT table_name, column_name, data_type FROM schema_snapshots
WHERE snapshot_date = (SELECT MAX(snapshot_date) FROM schema_snapshots)
)
SELECT 'added' AS change, c.* FROM current c
LEFT JOIN previous p USING (table_name, column_name) WHERE p.column_name IS NULL
UNION ALL
SELECT 'dropped', p.* FROM previous p
LEFT JOIN current c USING (table_name, column_name) WHERE c.column_name IS NULL;Expected output: an empty result on normal runs, or a precise list of added and dropped columns when drift happened.
- Also catch type changes, which the added/dropped diff can miss when a column keeps its name:
SELECT p.table_name, p.column_name, p.data_type AS old_type, c.data_type AS new_type
FROM previous p JOIN current c USING (table_name, column_name)
WHERE p.data_type != c.data_type;Expected output: columns whose types changed, like integer to bigint or varchar(50) to text. Type changes break loads just as thoroughly as dropped columns.
- Set the policy per change type: fail hard on dropped or retyped columns, warn on added ones:
if dropped_or_retyped:
raise ValueError(f"schema drift detected: {changes}")
elif added:
log_warning(f"new columns: {changes}")Expected output: breaking changes stop the pipeline before bad data lands; additive changes get noticed without blocking everything.
- Store every snapshot with a date so you have a history of when drift happened:
INSERT INTO schema_snapshots
SELECT table_name, column_name, data_type, ordinal_position, CURRENT_DATE
FROM information_schema.columns WHERE table_schema = 'staging';Expected output: a dated history. When something breaks next month, you can pinpoint exactly which day the schema changed.
Variant phrasings
schema change detection data pipeline
The snapshot-and-diff pattern above. It is cheap, it runs in SQL, and it catches the changes that break loads: drops, renames (which look like drop plus add), and type changes.
upstream schema changed broke my pipeline
Add the check as the first task of the pipeline so the next change fails fast with a clear message instead of producing corrupt output halfway through the DAG.
dbt source schema drift
Use dbt source freshness plus a custom schema test, or run the snapshot query as a dbt model and test the diff. The concept is identical; only the harness changes.
Why it happens
Producers evolve their schemas as their applications change, and nothing in the typical pipeline forces them to announce it. Consumers, meanwhile, write SQL that assumes specific columns and types. The gap between those two realities is where silent corruption and loud failures both come from. Detecting drift doesnt prevent the change; it converts a mystery into a notification.
Edge cases
- Type widening (integer to bigint) is usually safe; consider warning instead of failing for known-safe widenings.
- Column reordering is noise in most databases; compare by name, not by ordinal position.
- Case-only changes matter in case-sensitive databases and are invisible in others; normalize case in the comparison if your DB folds it.
- Partition column changes can break readers even when the logical schema looks fine; include partition metadata in the snapshot if you use partitioned tables.
- Renames appear as a drop plus an add; if you need true rename detection, track it with producer changelogs rather than diffing.
Provenance
Resolved from the public thread: https://vectle.com/posts/pst_5dGSQ0g5oTmtARJ9gzw26w
Maintainer review
No maintainer verification is recorded for this version.
This records the version a maintainer checked. It does not assert that the version is the latest upstream release.