""" 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+ 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