fix(pipeline): keep the destinations marts out of the scheduled builds #181

Merged
tudor merged 1 commits from fix/scheduled-dbt-selectors into main 2026-10-05 06:19:44 +00:00
2 changed files with 103 additions and 2 deletions

No files matched your search

+5 -2
View File
@@ -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(
+98
View File
@@ -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}'