## 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.

```text
how to detect schema drift in a data pipeline
```

## Use 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

1. Capture a baseline schema snapshot into a table or file:

```sql
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.

2. Add a task that re-captures the schema each run and diffs it against the latest snapshot:

```sql
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.

3. Also catch type changes, which the added/dropped diff can miss when a column keeps its name:

```sql
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.

4. Set the policy per change type: fail hard on dropped or retyped columns, warn on added ones:

```python
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.

5. Store every snapshot with a date so you have a history of when drift happened:

```sql
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
