Handling Schema Drift in Interconnection Queue Exports
The scenario: a monthly interconnection-queue load runs clean, the row count is normal, and the total
queued capacity has fallen by 94 percent. The ISO renamed capacity_mw to mw_capacity in the July
export; the loader’s get("capacity_mw", 0) did exactly what it was told, and every project in the
file now has zero megawatts. This page builds the ingestion that fails instead, and it is the schema
half of
geospatial data ingestion pipelines.
Root-cause analysis
Schema drift in public energy data is routine, and four shapes cover almost all of it.
- A renamed column. The most common and the most damaging, because a defaulted lookup turns a
missing column into a plausible value rather than an error.
df.get(col, 0)anddict.get(key, "")are the two lines that convert a loud failure into a silent one. - A changed unit. Capacity published in kilowatts where it was megawatts, or voltage in volts where it was kilovolts. The column name is unchanged, the type is unchanged, and every value is a thousand times off — which passes a range check that was written generously.
- A new nesting level. A flat CSV becomes a JSON payload with the records under a
datakey, or a single-value field becomes a list. The parser either raises immediately or, worse, reads the first element and drops the rest. - A widened domain. A status field gains a new value — “withdrawn-pending” alongside “withdrawn” — and any filter written as an exact match silently reclassifies those projects.
Pre-flight validation: a drift report, not a boolean
The useful pre-flight compares the incoming file against the contract and reports every difference, rather than answering yes or no. A load that fails needs to say what changed.
from dataclasses import dataclass, field
import pandas as pd
@dataclass
class SchemaContract:
version: str
required: dict[str, str] # column -> pandas dtype string
optional: dict[str, str] = field(default_factory=dict)
units: dict[str, str] = field(default_factory=dict)
domains: dict[str, set] = field(default_factory=dict)
def drift_report(df: pd.DataFrame, contract: SchemaContract) -> dict:
"""Every difference between the incoming frame and the contract, named."""
incoming = set(df.columns)
expected = set(contract.required) | set(contract.optional)
missing = sorted(set(contract.required) - incoming)
added = sorted(incoming - expected)
retyped = {
c: (str(df[c].dtype), contract.required[c])
for c in contract.required
if c in incoming and str(df[c].dtype) != contract.required[c]
}
domain_breaks = {}
for col, allowed in contract.domains.items():
if col in incoming:
unseen = sorted(set(df[col].dropna().unique()) - allowed)
if unseen:
domain_breaks[col] = unseen[:10]
# A renamed column usually appears as one missing and one added with a similar name.
likely_renames = [
(m, a) for m in missing for a in added
if _similar(m, a)
]
return {
"contract_version": contract.version,
"missing_required": missing,
"unexpected_columns": added,
"type_changes": retyped,
"domain_breaks": domain_breaks,
"likely_renames": likely_renames,
"clean": not (missing or retyped or domain_breaks),
}
def _similar(a: str, b: str) -> bool:
"""Cheap rename heuristic: same tokens, different order or separator."""
norm = lambda s: sorted(s.lower().replace("-", "_").split("_"))
return norm(a) == norm(b)
The likely_renames heuristic is worth the twelve lines: it turns “capacity_mw is missing and
mw_capacity appeared” into a one-line suggestion, which is the difference between a five-minute fix
and an afternoon of comparing files.
Fix implementation
The loader below validates against a versioned contract, quarantines rather than coerces, and treats a unit check as a first-class assertion rather than a range check.
import pandas as pd
QUEUE_V3 = SchemaContract(
version="queue.v3",
required={
"queue_id": "object",
"poi_name": "object",
"capacity_mw": "float64",
"voltage_kv": "float64",
"status": "object",
"state": "object",
},
optional={"withdrawn_date": "datetime64[ns]", "operator": "object"},
units={"capacity_mw": "MW", "voltage_kv": "kV"},
domains={"status": {"active", "withdrawn", "in-service", "suspended"}},
)
def load_queue_export(path: str, *, contract: SchemaContract = QUEUE_V3) -> tuple[pd.DataFrame, dict]:
"""Load an export, or fail with a report that names what changed."""
raw = pd.read_csv(path)
report = drift_report(raw, contract)
if report["missing_required"]:
raise ValueError(
f"{contract.version}: missing {report['missing_required']}; "
f"likely renames {report['likely_renames'] or 'none detected'}"
)
if report["type_changes"]:
raise ValueError(f"{contract.version}: type changes {report['type_changes']}")
# Unit sanity: a magnitude check, not a range check. Queue capacities live in
# the tens to hundreds of MW; a kW export lands three orders of magnitude out.
median_mw = float(raw["capacity_mw"].median())
if median_mw > 5_000:
raise ValueError(
f"median capacity {median_mw:,.0f} suggests kW, not MW — unit drift in {path}"
)
median_kv = float(raw["voltage_kv"].median())
if median_kv > 2_000:
raise ValueError(f"median voltage {median_kv:,.0f} suggests volts, not kV — unit drift")
unknown = report["domain_breaks"].get("status", [])
quarantine = raw[raw["status"].isin(unknown)] if unknown else raw.iloc[0:0]
clean = raw.drop(index=quarantine.index)
return clean, {**report, "quarantined_rows": len(quarantine), "loaded_rows": len(clean)}
Fallback routing and performance tuning
- Version the contract, do not edit it. When a source genuinely changes, add
queue.v4and keepv3; the loader then reports which version an old file matches, and historical reloads still work. - Quarantine unknown domain values, do not drop them. A new status value is information about the source, and dropping the rows makes it invisible while changing every total.
- Alert on the delta, not the level. A quarantine rate that jumps from 0.2 to 6 percent overnight is the signal; the absolute level says more about the source’s habits than about this run.
- Keep the raw payload. Reprocessing a month after fixing a contract is a minute if the bytes were kept and a re-download if they were not — and public portals do not always serve history.
- Never coerce silently.
errors="coerce"turns unparseable values into NaN, which then fails a nullability check and reports the wrong cause. Parse strictly and route failures to quarantine.
Downstream validation
def assert_load_is_comparable(current: dict, previous: dict, *, max_shift: float = 0.25) -> None:
"""Compare this load against the last one — drift usually shows up as a step change."""
for field_ in ("loaded_rows", "total_capacity_mw"):
prev, now = previous.get(field_), current.get(field_)
if not prev:
continue
shift = abs(now - prev) / prev
assert shift <= max_shift, (
f"{field_} moved {shift:.0%} against the previous load "
f"({prev:,.0f} → {now:,.0f}) — inspect before publishing"
)
prev_rate = previous.get("quarantined_rows", 0) / max(previous.get("loaded_rows", 1), 1)
now_rate = current.get("quarantined_rows", 0) / max(current.get("loaded_rows", 1), 1)
assert now_rate - prev_rate < 0.05, (
f"quarantine rate rose from {prev_rate:.1%} to {now_rate:.1%} — schema drift is likely"
)
Detecting drift before the load, from the header alone
Most drift is visible in the first kilobyte of a file, which means it can be caught before a multi-gigabyte download completes or a partition is replaced.
For a CSV, reading the header row and the first hundred data rows is enough to run the entire drift report: column names, inferred types, the domain of any low-cardinality field, and the median magnitude of the numeric columns. Ninety-nine percent of the file adds nothing to that judgement. For a JSON payload, the equivalent is the top-level keys and the first record.
That matters operationally because it changes where the failure lands. A drift check that runs after the download and before the write turns “the nightly job failed at 04:20 having written half a partition” into “the nightly job refused the July export at 03:02 and left last month’s data live”. The second is a triage task the next morning; the first is an incident.
The same sampling makes a dry-run cheap. Fetching only the headers of every partition due tonight — 51 states, a kilobyte each — takes seconds and reports exactly which sources drifted, which is enough to decide whether the run should proceed at all. Where the source supports a range request or a schema endpoint, the check costs no bytes worth counting.
One caveat worth stating: a sampled domain check is not exhaustive. A new status value that appears in row 90,000 will not be in the first hundred rows, so the sampled check catches the shapes that appear in the header and the full check still runs after the load. Sampling is an early warning, not a replacement.
Frequently asked questions
Should the loader try to auto-correct a detected rename?
No — report it and stop. An automatic rename is a guess about semantics made by a heuristic that only
compared strings, and the failure mode is a column mapped to the wrong meaning with no record that a
decision was made. Reporting capacity_mw missing and mw_capacity present takes seconds to
confirm, and the fix belongs in a versioned contract where it is visible in a diff.
How do I catch a unit change that stays within a plausible range?
By comparing against the previous load rather than against an absolute range. A capacity column that moves from a median of 120 to a median of 0.12 is obviously wrong; one that moves from 120 to 132 is probably real growth. The step-change assertion above catches the first class, and no automated check reliably catches the second — which is why the previous-load comparison is worth more than a wider range check.
What belongs in the contract versus in the validation schema?
The contract describes the file as delivered: column names, types, units and domains. The validation schema — pandera or equivalent — describes the record as the pipeline needs it: ranges, nullability, cross-field consistency. Keeping them separate means a source change updates the contract and a business-rule change updates the schema, and neither edit touches the other.
How many contract versions should be kept?
All of them, because they are a few dozen lines each and they are what makes historical reprocessing possible. Tag each stored raw payload with the contract version it matched at load time, and a reprocess two years later resolves the right parser without anyone remembering which month the format changed.
Does this apply to spatial columns too?
Yes, and the geometry equivalents are worth naming: a CRS that changes between exports, a geometry column that arrives as WKT where it was WKB, and coordinates that swap axis order. All three are detectable with the same magnitude reasoning — a longitude column whose median is 35 rather than −101 has been transposed, and that check costs one line.
What should the alert actually say?
The contract version, the file, the specific differences, and the previous load’s figures for comparison. An alert that says “schema validation failed” starts an investigation; one that says “queue.v3: missing capacity_mw, likely renamed to mw_capacity; 1,284 rows unloaded; previous load 1,247 rows / 42,180 MW” ends it.
Related
- Geospatial Data Ingestion Pipelines — the parent ingestion contract
- Incremental, Idempotent Loading of Grid Datasets with Hive Partitions — writing what survives this validation
- Enforcing Voltage Class Schemas with pandera — the record-level schema this contract feeds
- Open Energy Data Portals — where these exports come from and how often they change