Compare commits

..
Author SHA1 Message Date
TudorandClaude Opus 5.5 65a2619e1d fix(pipeline): keep the destinations marts out of the scheduled builds
PR Checks / Frontend Typecheck + Tests (pull_request) Successful in 1m18s
PR Checks / Backend Smoke (pull_request) Successful in 11s
PR Checks / Build Backend (no push) (pull_request) Successful in 18s
PR Checks / Build Frontend (no push) (pull_request) Successful in 1m35s
PR Checks / Build Pipeline (no push) (pull_request) Successful in 1m21s
PR Checks / AI Code Review (Claude) (pull_request) Successful in 19s
fact_ks4_destinations and fact_ks5_destinations join dim_school, so the
daily build's stg_gias_establishments+ and the monthly Ofsted build's
dim_school+ both selected them. They also read stg_ees_ks4/ks5_destinations,
which only the manually triggered EES DAG builds. Where that DAG hasn't run
since the destinations models landed, dbt_build fails with "relation
staging.stg_ees_ks4_destinations does not exist", and sync_typesense and
invalidate_cache never run. Production's register data has been stuck at
about 25 Aug 2026.

Both builds now exclude the descendants of the two EES staging models, as
the daily build already does for the KS2/KS4 lineage models. The EES DAG
still rebuilds the marts when their data changes.

test_dag_selectors reads the model graph from the SQL (CI has no dbt) and
checks that every scheduled build only reads models it or the daily build
builds. It failed for the daily and monthly Ofsted builds before this change.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-03 22:21:27 +01:00
5 changed files with 121 additions and 98 deletions

No files matched your search

+5 -2
View File
@@ -106,9 +106,12 @@ print(f'Validation passed: {{count}} GIAS rows')
""",
)
# Marts fed by annual EES staging models are rebuilt by the EES DAG, even
# when they join dim_school. Selecting them here fails in any database
# where that DAG hasn't run (pipeline/tests/test_dag_selectors.py).
dbt_build = BashOperator(
task_id="dbt_build",
bash_command=f"cd {PIPELINE_DIR}/transform && {DBT_BIN} build --profiles-dir . --target production --select stg_gias_establishments+ stg_gias_links+ gias_code_names+ --exclude int_ks2_with_lineage+ int_ks4_with_lineage+",
bash_command=f"cd {PIPELINE_DIR}/transform && {DBT_BIN} build --profiles-dir . --target production --select stg_gias_establishments+ stg_gias_links+ gias_code_names+ --exclude int_ks2_with_lineage+ int_ks4_with_lineage+ stg_ees_ks4_destinations+ stg_ees_ks5_destinations+",
)
sync_typesense = BashOperator(
@@ -143,7 +146,7 @@ with DAG(
dbt_build_ofsted = BashOperator(
task_id="dbt_build",
bash_command=f"cd {PIPELINE_DIR}/transform && {DBT_BIN} build --profiles-dir . --target production --select stg_ofsted_inspections+ int_ofsted_latest+ fact_ofsted_inspection+ dim_school+",
bash_command=f"cd {PIPELINE_DIR}/transform && {DBT_BIN} build --profiles-dir . --target production --select stg_ofsted_inspections+ int_ofsted_latest+ fact_ofsted_inspection+ dim_school+ --exclude stg_ees_ks4_destinations+ stg_ees_ks5_destinations+",
)
sync_typesense_ofsted = BashOperator(
@@ -1,26 +0,0 @@
"""Read a GIAS extract from the raw bytes of the download.
GIAS writes its CSVs in Windows-1252 and sends no charset, so `resp.text`
leaves requests to guess the codec. On 3 Oct 2026 it guessed windows-1250 and
"à" became "ŕ". Decode the bytes ourselves instead.
"""
from __future__ import annotations
import io
import pandas as pd
GIAS_ENCODING = "cp1252"
def read_gias_csv(content: bytes, logger=None) -> pd.DataFrame:
"""Every column as a string; a blank cell stays ''."""
# Windows-1252 leaves five bytes undefined. One stray byte must not stop
# the daily refresh of every school, so it becomes U+FFFD and is logged.
text = content.decode(GIAS_ENCODING, errors="replace")
undecodable = text.count("�")
if undecodable and logger is not None:
logger.warning("%d byte(s) in the GIAS extract could not be decoded as %s",
undecodable, GIAS_ENCODING)
return pd.read_csv(io.StringIO(text), dtype=str, keep_default_na=False)
@@ -7,8 +7,6 @@ from datetime import date, timedelta
from singer_sdk import Stream, Tap
from singer_sdk import typing as th
from tap_uk_gias.gias_csv import read_gias_csv
GIAS_URL_TEMPLATE = (
"https://ea-edubase-api-prod.azurewebsites.net"
"/edubase/downloads/public/edubasealldata{date}.csv"
@@ -76,6 +74,9 @@ class GIASEstablishmentsStream(Stream):
def get_records(self, context):
"""Download GIAS CSV and yield rows."""
import io
import pandas as pd
import requests
today = date.today()
@@ -93,7 +94,12 @@ class GIASEstablishmentsStream(Stream):
resp.raise_for_status()
df = read_gias_csv(resp.content, self.logger)
df = pd.read_csv(
io.StringIO(resp.text),
encoding="latin-1",
dtype=str,
keep_default_na=False,
)
for _, row in df.iterrows():
record = row.to_dict()
@@ -120,6 +126,9 @@ class GIASLinksStream(Stream):
def get_records(self, context):
"""Download GIAS links CSV and yield rows."""
import io
import pandas as pd
import requests
today = date.today()
@@ -137,7 +146,12 @@ class GIASLinksStream(Stream):
resp.raise_for_status()
df = read_gias_csv(resp.content, self.logger)
df = pd.read_csv(
io.StringIO(resp.text),
encoding="latin-1",
dtype=str,
keep_default_na=False,
)
for _, row in df.iterrows():
record = row.to_dict()
+98
View File
@@ -0,0 +1,98 @@
"""Every scheduled dbt build must only build models whose parents exist.
The daily GIAS build selects `stg_gias_establishments+`, so any mart that joins
dim_school joins the daily build too. When such a mart also reads a staging
model that only the manually triggered EES DAG builds, the daily build fails in
any database where that DAG has not run since. Sync and cache invalidation then
never run either. The destinations marts did this from late August 2026.
The graph is read from the model SQL, because CI has no dbt.
"""
import re
from collections import defaultdict
from pathlib import Path
import pytest
PIPELINE = Path(__file__).resolve().parents[1]
MODELS = PIPELINE / 'transform' / 'models'
DAG_FILE = PIPELINE / 'dags' / 'school_data_pipeline.py'
REF = re.compile(r"ref\(\s*'([a-z0-9_]+)'\s*\)")
DBT_BUILD = re.compile(r'dbt_build\w*\s*=\s*BashOperator\(.*?build --profiles-dir \. --target production ([^"]+)"', re.S)
DAG_ID = re.compile(r'dag_id="([a-z0-9_]+)"')
DAILY = 'school_data_daily'
# dim_school reads int_ofsted_latest only when the relation exists
# (adapter.get_relation), so a missing table is not a failure.
OPTIONAL_PARENTS = {'int_ofsted_latest'}
def model_parents():
"""{model: models it refs}. Seeds are left out: they are loaded once and always exist."""
sql = {p.stem: p.read_text() for p in MODELS.rglob('*.sql')}
return {name: set(REF.findall(text)) & set(sql) for name, text in sql.items()}
def downstream(node, children):
seen, stack = {node}, [node]
while stack:
for child in children[stack.pop()]:
if child not in seen:
seen.add(child)
stack.append(child)
return seen
def expand(tokens, children):
out = set()
for token in tokens:
out |= downstream(token[:-1], children) if token.endswith('+') else {token}
return out
def scheduled_builds():
"""{dag_id: dbt selection arguments} for every dbt build in the DAG file."""
text = DAG_FILE.read_text()
starts = [(m.start(), m.group(1)) for m in DAG_ID.finditer(text)]
builds = {}
for i, (start, dag_id) in enumerate(starts):
end = starts[i + 1][0] if i + 1 < len(starts) else len(text)
found = DBT_BUILD.search(text, start, end)
if found:
builds[dag_id] = found.group(1)
return builds
def selected_models(args, parents):
children = defaultdict(set)
for model, ps in parents.items():
for p in ps:
children[p].add(model)
select = re.search(r'--select (.+?)(?= --exclude|$)', args).group(1).split()
excluded = re.search(r'--exclude (.+)$', args)
exclude = excluded.group(1).split() if excluded else []
return (expand(select, children) - expand(exclude, children)) & set(parents)
PARENTS = model_parents()
BUILDS = scheduled_builds()
DAILY_MODELS = selected_models(BUILDS[DAILY], PARENTS)
def test_every_dag_with_a_dbt_build_is_parsed():
assert set(BUILDS) == {
'school_data_daily', 'school_data_monthly_ofsted', 'school_data_annual_ees',
'school_data_annual_idaci', 'school_data_annual_distance',
}
@pytest.mark.parametrize('dag_id', sorted(BUILDS))
def test_selected_models_only_read_models_that_exist(dag_id):
selected = selected_models(BUILDS[dag_id], PARENTS)
# The daily build is the base layer: other DAGs may rely on what it builds.
available = selected | OPTIONAL_PARENTS | (DAILY_MODELS if dag_id != DAILY else set())
missing = {model: sorted(PARENTS[model] - available) for model in sorted(selected)
if PARENTS[model] - available}
assert missing == {}, f'{dag_id} builds models whose parents it never builds: {missing}'
-66
View File
@@ -1,66 +0,0 @@
"""GIAS publishes its extracts in Windows-1252 and declares no charset.
The tap used to hand pandas `resp.text`, so requests guessed the codec.
On 3 Oct 2026 it guessed windows-1250, and "St Thomas à Becket" was stored
as "St Thomas ŕ Becket". The `encoding=` passed to read_csv did nothing,
because the text was already decoded.
"""
import importlib.util
import logging
from pathlib import Path
import pytest
MODULE = (Path(__file__).resolve().parents[1] / 'plugins' / 'extractors' / 'tap-uk-gias'
/ 'tap_uk_gias' / 'gias_csv.py')
@pytest.fixture
def gias_csv():
spec = importlib.util.spec_from_file_location('gias_csv', MODULE)
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
return module
# Byte for byte as GIAS writes it: 0xE0 à, 0x92 ’, 0xE9 é, 0xB0 °, 0xE7 ç.
EXTRACT = (
b'"URN","EstablishmentName","HeadLastName"\r\n'
b'"138950","St Thomas \xe0 Becket Catholic Secondary School","Smith"\r\n'
b'"100000","The Dean and Chapter of St Paul\x92s Cathedral","Pr\xe9vert"\r\n'
b'"140677","North Star 180\xb0","Fran\xe7ois"\r\n'
b'"100001","No head recorded",""\r\n'
)
def test_names_decode_as_windows_1252(gias_csv):
df = gias_csv.read_gias_csv(EXTRACT)
assert list(df['EstablishmentName']) == [
'St Thomas à Becket Catholic Secondary School',
'The Dean and Chapter of St Paul’s Cathedral',
'North Star 180°',
'No head recorded',
]
assert list(df['HeadLastName']) == ['Smith', 'Prévert', 'François', '']
def test_the_codec_requests_guessed_is_not_used(gias_csv):
# What the tap stored on 3 Oct: the same bytes read as windows-1250.
assert 'ŕ' in EXTRACT.decode('cp1250')
names = ' '.join(gias_csv.read_gias_csv(EXTRACT)['EstablishmentName'])
assert 'ŕ' not in names
def test_values_stay_strings(gias_csv):
df = gias_csv.read_gias_csv(EXTRACT)
assert df.loc[0, 'URN'] == '138950'
def test_a_byte_windows_1252_leaves_undefined_does_not_stop_the_load(gias_csv, caplog):
# 0x81 has no Windows-1252 character. One odd name must not block the daily
# refresh of every school, but it must be visible in the log.
extract = b'"URN","EstablishmentName"\r\n"100002","Odd \x81 Name"\r\n'
with caplog.at_level(logging.WARNING):
df = gias_csv.read_gias_csv(extract, logger=logging.getLogger('gias'))
assert df.loc[0, 'EstablishmentName'] == 'Odd � Name'
assert 'could not be decoded' in caplog.text