diff --git a/pipeline/dags/school_data_pipeline.py b/pipeline/dags/school_data_pipeline.py index a8b5a2a..2b27af7 100644 --- a/pipeline/dags/school_data_pipeline.py +++ b/pipeline/dags/school_data_pipeline.py @@ -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( diff --git a/pipeline/tests/test_dag_selectors.py b/pipeline/tests/test_dag_selectors.py new file mode 100644 index 0000000..c41a55e --- /dev/null +++ b/pipeline/tests/test_dag_selectors.py @@ -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}'