Files
school_compare/pipeline/dags/school_data_pipeline.py
T
TudorandClaude Opus 5.5 65a2619e1d
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
fix(pipeline): keep the destinations marts out of the scheduled builds
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

305 lines
11 KiB
Python

"""
School Data Pipeline — Airflow DAG
Orchestrates the full ELT pipeline:
Extract (Meltano) → Validate → Transform (dbt) → Geocode → Sync Typesense → Invalidate Cache
Schedule:
- GIAS: Daily at 03:00
- Ofsted: 1st of month at 02:00
- EES datasets: Annual (triggered manually or on detected release)
"""
from __future__ import annotations
from datetime import datetime, timedelta
from airflow import DAG
from airflow.utils.task_group import TaskGroup
try:
from airflow.providers.standard.operators.bash import BashOperator
except ImportError:
from airflow.operators.bash import BashOperator
PIPELINE_DIR = "/opt/pipeline"
MELTANO_BIN = "meltano"
# Invoke dbt via the Python module rather than a bare `dbt` on PATH. The
# standalone dbt Fusion binary (dbt-core 2.x) shadows the pip-installed
# classic engine on some images and rejects the Postgres adapter
# (dbt1005). The module form always resolves to dbt-postgres ~=1.10.
DBT_BIN = "python -m dbt.cli.main"
default_args = {
"owner": "school-compare",
"depends_on_past": False,
"email_on_failure": False,
"retries": 1,
"retry_delay": timedelta(minutes=5),
}
# The backend caches the marts DataFrame at startup; after any rebuild the
# cache must be invalidated or the API serves stale (or empty) data until the
# container restarts.
INVALIDATE_CACHE_CMD = """
set -e
BACKEND_URL="${BACKEND_URL:-http://backend:80}"
ADMIN_KEY="${ADMIN_API_KEY:-changeme}"
echo "Calling $BACKEND_URL/api/admin/reload ..."
response=$(curl -s -o /tmp/reload_response.json -w "%{http_code}" \\
--connect-timeout 10 --max-time 120 \\
-X POST "$BACKEND_URL/api/admin/reload" \\
-H "X-API-Key: $ADMIN_KEY" \\
-H "Content-Type: application/json")
echo "HTTP status: $response"
cat /tmp/reload_response.json
if [ "$response" != "200" ]; then
echo "ERROR: backend cache reload failed (HTTP $response)"
exit 1
fi
"""
# ── Daily DAG (GIAS + downstream) ──────────────────────────────────────
with DAG(
dag_id="school_data_daily",
default_args=default_args,
description="Daily school data pipeline (GIAS extract → full transform)",
schedule="0 3 * * *",
start_date=datetime(2025, 1, 1),
catchup=False,
tags=["school-compare", "daily"],
) as daily_dag:
with TaskGroup("extract") as extract_group:
extract_gias = BashOperator(
task_id="extract_gias",
bash_command=f"cd {PIPELINE_DIR} && {MELTANO_BIN} run tap-uk-gias target-postgres",
)
validate_raw = BashOperator(
task_id="validate_raw",
bash_command=f"""
cd {PIPELINE_DIR} && python -c "
import psycopg2, os, sys
conn = psycopg2.connect(
host=os.environ.get('PG_HOST', 'localhost'),
port=os.environ.get('PG_PORT', '5432'),
user=os.environ.get('PG_USER', 'postgres'),
password=os.environ.get('PG_PASSWORD', 'postgres'),
dbname=os.environ.get('PG_DATABASE', 'school_compare'),
)
cur = conn.cursor()
cur.execute('SELECT count(*) FROM raw.gias_establishments')
count = cur.fetchone()[0]
conn.close()
if count < 20000:
print(f'WARN: GIAS only has {{count}} rows, expected 60k+', file=sys.stderr)
sys.exit(1)
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+ stg_ees_ks4_destinations+ stg_ees_ks5_destinations+",
)
sync_typesense = BashOperator(
task_id="sync_typesense",
bash_command=f"cd {PIPELINE_DIR} && python scripts/sync_typesense.py",
)
invalidate_cache = BashOperator(
task_id="invalidate_cache",
bash_command=INVALIDATE_CACHE_CMD,
)
extract_group >> validate_raw >> dbt_build >> sync_typesense >> invalidate_cache
# ── Monthly DAG (Ofsted) ───────────────────────────────────────────────
with DAG(
dag_id="school_data_monthly_ofsted",
default_args=default_args,
description="Monthly Ofsted MI extraction and transform",
schedule="0 2 1 * *",
start_date=datetime(2025, 1, 1),
catchup=False,
tags=["school-compare", "monthly"],
) as monthly_ofsted_dag:
extract_ofsted = BashOperator(
task_id="extract_ofsted",
bash_command=f"cd {PIPELINE_DIR} && {MELTANO_BIN} run tap-uk-ofsted target-postgres",
)
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+ --exclude stg_ees_ks4_destinations+ stg_ees_ks5_destinations+",
)
sync_typesense_ofsted = BashOperator(
task_id="sync_typesense",
bash_command=f"cd {PIPELINE_DIR} && python scripts/sync_typesense.py",
)
invalidate_cache_ofsted = BashOperator(
task_id="invalidate_cache",
bash_command=INVALIDATE_CACHE_CMD,
)
extract_ofsted >> dbt_build_ofsted >> sync_typesense_ofsted >> invalidate_cache_ofsted
# ── Annual DAG (EES: KS2, KS4, Census, Admissions) ───────────────────
with DAG(
dag_id="school_data_annual_ees",
default_args=default_args,
description="Annual EES data extraction (KS2, KS4, Census, Admissions)",
schedule=None, # Triggered manually when new releases are published
start_date=datetime(2025, 1, 1),
catchup=False,
tags=["school-compare", "annual"],
) as annual_ees_dag:
with TaskGroup("extract_ees") as extract_ees_group:
# Runs all tap-uk-ees streams, including the new ees_ks2_national stream
extract_ees = BashOperator(
task_id="extract_ees",
bash_command=f"cd {PIPELINE_DIR} && {MELTANO_BIN} run tap-uk-ees target-postgres",
)
# Destinations come from the EES query API rather than a release ZIP,
# so they are a separate tap. Run after extract_ees rather than beside
# it: both write to the raw schema and the loader is happier serial.
extract_ees_destinations = BashOperator(
task_id="extract_ees_destinations",
bash_command=f"cd {PIPELINE_DIR} && {MELTANO_BIN} run tap-uk-ees-destinations target-postgres",
)
extract_ees >> extract_ees_destinations
dbt_build_ees = BashOperator(
task_id="dbt_build",
bash_command=f"cd {PIPELINE_DIR}/transform && {DBT_BIN} build --profiles-dir . --target production --select stg_ees_ks2+ stg_legacy_ks2+ stg_ees_ks4+ stg_legacy_ks4+ stg_ees_census+ stg_ees_admissions+ stg_ees_ks2_national+ stg_ees_ks4_national+ stg_ees_ks4_destinations+ stg_ees_ks5_destinations+ stg_ees_ks4_destinations_national+ stg_ees_ks5_destinations_national+",
)
sync_typesense_ees = BashOperator(
task_id="sync_typesense",
bash_command=f"cd {PIPELINE_DIR} && python scripts/sync_typesense.py",
)
invalidate_cache_ees = BashOperator(
task_id="invalidate_cache",
bash_command=INVALIDATE_CACHE_CMD,
)
extract_ees_group >> dbt_build_ees >> sync_typesense_ees >> invalidate_cache_ees
# ── Annual DAG (IDACI Deprivation) ────────────────────────────────────
with DAG(
dag_id="school_data_annual_idaci",
default_args=default_args,
description="Annual IDACI deprivation index extraction and transform",
schedule=None, # Triggered manually when new IDACI release is published
start_date=datetime(2025, 1, 1),
catchup=False,
tags=["school-compare", "annual"],
) as annual_idaci_dag:
extract_idaci = BashOperator(
task_id="extract_idaci",
bash_command=f"cd {PIPELINE_DIR} && {MELTANO_BIN} run tap-uk-idaci target-postgres",
)
dbt_build_idaci = BashOperator(
task_id="dbt_build",
bash_command=f"cd {PIPELINE_DIR}/transform && {DBT_BIN} build --profiles-dir . --target production --select stg_idaci+ fact_deprivation+",
)
invalidate_cache_idaci = BashOperator(
task_id="invalidate_cache",
bash_command=INVALIDATE_CACHE_CMD,
)
extract_idaci >> dbt_build_idaci >> invalidate_cache_idaci
# ── Annual DAG (Last distance offered) ────────────────────────────────
with DAG(
dag_id="school_data_annual_distance",
default_args=default_args,
description="Last distance offered (LA admission cut-offs) extraction and transform",
# Councils publish on allocation day (March for secondary, April for
# primary) and each one on its own timetable, so there is no date worth
# scheduling against. Triggered manually after a collection run refreshes
# the CSV.
schedule=None,
start_date=datetime(2025, 1, 1),
catchup=False,
tags=["school-compare", "annual"],
) as annual_distance_dag:
extract_distance = BashOperator(
task_id="extract_distance",
bash_command=f"cd {PIPELINE_DIR} && {MELTANO_BIN} run tap-uk-school-distance target-postgres",
)
# Coverage is the thing that silently rots here: the CSV is assembled by
# hand from council publications, so a collection run that half-failed
# produces a valid file with a fraction of the schools in it. A row count
# alone would not catch that — losing an entire local authority leaves the
# total looking healthy — so the floor is checked on distinct LAs too.
validate_distance = BashOperator(
task_id="validate_distance",
bash_command=f"""
cd {PIPELINE_DIR} && python -c "
import psycopg2, os, sys
conn = psycopg2.connect(
host=os.environ.get('PG_HOST', 'localhost'),
port=os.environ.get('PG_PORT', '5432'),
user=os.environ.get('PG_USER', 'postgres'),
password=os.environ.get('PG_PASSWORD', 'postgres'),
dbname=os.environ.get('PG_DATABASE', 'school_compare'),
)
cur = conn.cursor()
cur.execute('SELECT count(*), count(distinct urn), count(distinct la_code) FROM raw.school_distance_offered')
rows, urns, las = cur.fetchone()
conn.close()
print(f'Loaded {{rows}} rows, {{urns}} schools, {{las}} local authorities')
if rows < 7000 or urns < 3000 or las < 45:
print('ERROR: distance extract is short of expected coverage '
'(baseline 2026-08: 9128 rows / 3726 schools / 57 LAs)', file=sys.stderr)
sys.exit(1)
print('Validation passed')
"
""",
)
dbt_build_distance = BashOperator(
task_id="dbt_build",
bash_command=f"cd {PIPELINE_DIR}/transform && {DBT_BIN} build --profiles-dir . --target production --select stg_school_distance+ fact_admission_distance+",
)
invalidate_cache_distance = BashOperator(
task_id="invalidate_cache",
bash_command=INVALIDATE_CACHE_CMD,
)
extract_distance >> validate_distance >> dbt_build_distance >> invalidate_cache_distance