Files
school_compare/pipeline/dags/school_data_pipeline.py
T
TudorandClaude Opus 5 88c653215d
PR Checks / Frontend Typecheck + Tests (pull_request) Successful in 1m3s
PR Checks / Backend Smoke (pull_request) Successful in 7s
PR Checks / Build Backend (no push) (pull_request) Successful in 31s
PR Checks / Build Frontend (no push) (pull_request) Successful in 44s
PR Checks / Build Pipeline (no push) (pull_request) Successful in 1m12s
PR Checks / AI Code Review (Claude) (pull_request) Failing after 2m7s
feat(admissions): show the last distance offered where councils publish it
Adds the cut-off distance a parent actually asks about — "how close do we
need to live?" — end to end: a Singer tap, dbt staging and mart models, an
Airflow DAG, and a tile on both detail templates. 3,597 schools across 57
local authorities carry a figure; the rest are unchanged.

There is no national source for this. Each LA publishes its own cut-offs in
its own format, and the collected CSV is transcribed from PDFs, spreadsheets
and web pages — so most of the work here is deciding what is safe to show.

Data
  * tap-uk-school-distance loads the CSV verbatim into raw. Keyed on
    (urn, year, school_name), because school_name carries the admission
    route: (urn, year) alone collides on 118 keys and a reload would have
    silently dropped every band but one.
  * stg_school_distance applies a 25 m – 25 km plausibility band. The source
    contains 0.0-mile rows (published where a school filled on a higher
    criterion), 1-metre cut-offs, and one reading 533 miles — ~4% of rows,
    all of which would put a visibly wrong number on a live page.
  * fact_admission_distance collapses routes to one row per school per year
    using the furthest, and keeps route_count so the page can say the figure
    is the widest of several bands rather than the one for a given child.

Serving
  * Kept out of fact_admissions: that mart is EES-derived and near-complete
    for England, this one covers 57 LAs, and the two refresh independently.
  * Latest year only. Coverage is ragged — a school may have 2021 and 2026
    and nothing between — so a history array would invite a trend line drawn
    through gaps that are absences of publication, not of a cut-off.
  * The Admissions section now renders on either source. 3% of the schools
    that render have a cut-off and no EES admissions row, and gating on
    admissions alone would have hidden the figure on those pages.

Interface
  * The year travels with the figure everywhere it appears; a cut-off
    detached from its admissions round is not a fact about anything.
  * "Not a fixed catchment — it moves every year" sits under every instance,
    because that is the inference a parent will otherwise draw.
  * Replaces a hardcoded "Historical distance cut-off data is not available
    for this school" that appeared on every secondary page, including the
    ones whose council does publish it. The absence is now stated only when
    it is real, and names the authority that would hold it.

The tint costs the muted tokens their AA margin: measured on the composited
backdrop (not the computed one, which reports the untinted card), --text-muted
falls to 4.09:1 in dark theme. The tile uses --text-secondary instead — 6.50:1
dark, 6.60:1 light.

The DAG is manual, like the other annual ones: councils publish on allocation
day, each on its own timetable, so there is no date worth scheduling against.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01WDvkyqqHABm4bmth2kjAxE
2026-08-15 22:48:30 +01:00

292 lines
10 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')
"
""",
)
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+",
)
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+",
)
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",
)
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+",
)
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