EES writes 'c' where a figure is withheld and the categories sum to the cohort, so counts and percentages are emitted as text with the sentinel intact. safe_numeric must never be pointed at them. School rows and the England reference need different establishment pins: at national level selective schools, studios and UTCs are separate populations rather than labels, so leaving establishment open multiplies 30 rows into 190. Two queries per period, each keeping its own level. Verified against the live API for 2022/23: 135,240 school records over 4,508 schools, exactly 30 each, no duplicate keys, 31,382 sentinels kept. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01BvdDKvFFSZuMVDH5fEyTob
302 lines
11 KiB
Python
302 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')
|
|
"
|
|
""",
|
|
)
|
|
|
|
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",
|
|
)
|
|
|
|
# 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+",
|
|
)
|
|
|
|
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
|