From f15d5b9f3fee941c17a4648d65a56b5b4fb4cf14 Mon Sep 17 00:00:00 2001 From: BonelessWater Date: Thu, 8 Oct 2026 15:28:54 -0400 Subject: [PATCH] fix(data-ngin): fail the DAG run when the pipeline task fails Each DAG ended in a staleness_check task with trigger_rule=all_done. Airflow takes a run's state from its leaf tasks, so after run_pipeline failed the always-run leaf succeeded and the run was marked success. A manual run on the box hit exactly this: 2 of 29 Databento symbols failed, and the run showed green. Drop the warn-only task. The old production DAGs never had it, and the daily data_ngin.ops.check_data_freshness cron job covers it. A test now keeps all_done out of the DAG files. --- services/data-ngin/dags/data_pipeline_dag.py | 12 ++++---- .../data-ngin/dags/new_data_pipeline_dag.py | 12 ++++---- services/data-ngin/dags/tiingo_data_dag.py | 12 ++++---- .../data_ngin/application/pipeline_tasks.py | 29 ------------------- services/data-ngin/tests/test_prod_wiring.py | 14 +++++---- 5 files changed, 23 insertions(+), 56 deletions(-) diff --git a/services/data-ngin/dags/data_pipeline_dag.py b/services/data-ngin/dags/data_pipeline_dag.py index 6b5bd4cd..fb1c468a 100644 --- a/services/data-ngin/dags/data_pipeline_dag.py +++ b/services/data-ngin/dags/data_pipeline_dag.py @@ -31,12 +31,6 @@ def run_pipeline() -> None: pipeline_tasks.run_pipeline(CONFIG_NAME, run_type=conf.get("run_type", "scheduled")) -@task(trigger_rule="all_done") -def staleness_check() -> None: - # all_done: runs (and only warns) even when the pipeline task failed. - pipeline_tasks.check_staleness(CONFIG_NAME) - - @dag( dag_id="data_pipeline_dag", default_args=default_args, @@ -48,7 +42,11 @@ def staleness_check() -> None: max_active_runs=1, ) def data_pipeline_dag(): - run_pipeline() >> staleness_check() + # The pipeline task is the only task, so a failed run_* is a failed DAG + # run (and a GitHub issue via on_failure_callback). A trailing + # always-run task (a trigger rule that ignores upstream failures) would + # turn that into a green run. + run_pipeline() data_pipeline_dag() diff --git a/services/data-ngin/dags/new_data_pipeline_dag.py b/services/data-ngin/dags/new_data_pipeline_dag.py index 0ef3b469..3a87b1c9 100644 --- a/services/data-ngin/dags/new_data_pipeline_dag.py +++ b/services/data-ngin/dags/new_data_pipeline_dag.py @@ -31,12 +31,6 @@ def run_pipeline() -> None: pipeline_tasks.run_pipeline(CONFIG_NAME, run_type=conf.get("run_type", "scheduled")) -@task(trigger_rule="all_done") -def staleness_check() -> None: - # all_done: runs (and only warns) even when the pipeline task failed. - pipeline_tasks.check_staleness(CONFIG_NAME) - - @dag( dag_id="new_data_pipeline_dag", default_args=default_args, @@ -48,7 +42,11 @@ def staleness_check() -> None: max_active_runs=1, ) def new_data_pipeline_dag(): - run_pipeline() >> staleness_check() + # The pipeline task is the only task, so a failed run_* is a failed DAG + # run (and a GitHub issue via on_failure_callback). A trailing + # always-run task (a trigger rule that ignores upstream failures) would + # turn that into a green run. + run_pipeline() new_data_pipeline_dag() diff --git a/services/data-ngin/dags/tiingo_data_dag.py b/services/data-ngin/dags/tiingo_data_dag.py index 1bd7ca80..48a8469b 100644 --- a/services/data-ngin/dags/tiingo_data_dag.py +++ b/services/data-ngin/dags/tiingo_data_dag.py @@ -31,12 +31,6 @@ def run_tiingo_pipeline() -> None: pipeline_tasks.run_pipeline(CONFIG_NAME, run_type=conf.get("run_type", "scheduled")) -@task(trigger_rule="all_done") -def staleness_check() -> None: - # all_done: runs (and only warns) even when the pipeline task failed. - pipeline_tasks.check_staleness(CONFIG_NAME) - - @dag( dag_id="tiingo_data_dag", default_args=default_args, @@ -50,7 +44,11 @@ def staleness_check() -> None: max_active_runs=1, ) def tiingo_data_dag(): - run_tiingo_pipeline() >> staleness_check() + # The pipeline task is the only task, so a failed run_* is a failed DAG + # run (and a GitHub issue via on_failure_callback). A trailing + # always-run task (a trigger rule that ignores upstream failures) would + # turn that into a green run. + run_tiingo_pipeline() tiingo_data_dag() diff --git a/services/data-ngin/src/data_ngin/application/pipeline_tasks.py b/services/data-ngin/src/data_ngin/application/pipeline_tasks.py index 1a57890c..396ed024 100644 --- a/services/data-ngin/src/data_ngin/application/pipeline_tasks.py +++ b/services/data-ngin/src/data_ngin/application/pipeline_tasks.py @@ -10,7 +10,6 @@ import logging import os -from datetime import timedelta PACKAGE_CONFIG_DIR = os.path.join(os.path.dirname(os.path.dirname(__file__)), "config") @@ -45,31 +44,3 @@ def run_pipeline(config_name: str, run_type: str = "scheduled") -> None: orchestrator = Orchestrator(config=load_config(path)) asyncio.run(orchestrator.run()) logging.info(f"Pipeline {config_name} completed successfully.") - - -def check_staleness(config_name: str) -> None: - """ - Log a warning (never fail) when the pipeline's target table is older than - DATA_NGIN_STALENESS_THRESHOLD_DAYS (default 1). The cron'd - data_ngin.ops.check_data_freshness is the alerting path; this is the - in-Airflow signal next to the run that caused it. - """ - from data_ngin.domain.services import StalenessChecker - from data_ngin.infrastructure.repository.ohlcv_repository import OhlcvRepository - from data_ngin.utils.dynamic_loader import load_config - - repository = OhlcvRepository(config=load_config(config_path(config_name))) - repository.connect() - try: - latest_date = repository.get_latest_date() - finally: - repository.close() - - threshold_days = int(os.getenv("DATA_NGIN_STALENESS_THRESHOLD_DAYS", "1")) - report = StalenessChecker(max_staleness=timedelta(days=threshold_days)).check_staleness( - latest_date - ) - if report.is_stale: - logging.warning(report.message) - else: - logging.info(report.message) diff --git a/services/data-ngin/tests/test_prod_wiring.py b/services/data-ngin/tests/test_prod_wiring.py index c2cbaa6c..ef2e581f 100644 --- a/services/data-ngin/tests/test_prod_wiring.py +++ b/services/data-ngin/tests/test_prod_wiring.py @@ -3,6 +3,7 @@ import os import unittest +from pathlib import Path from unittest.mock import MagicMock, patch from data_ngin.application import pipeline_tasks @@ -138,12 +139,13 @@ async def _run() -> None: config = mock_orchestrator_cls.call_args.kwargs["config"] self.assertEqual(config["fetcher"]["class"], "TiingoFetcher") - @patch("data_ngin.infrastructure.repository.ohlcv_repository.OhlcvRepository") - def test_check_staleness_warns_without_raising(self, mock_repo_cls: MagicMock) -> None: - mock_repo_cls.return_value.get_latest_date.return_value = "2020-01-01" - with self.assertLogs(level="WARNING"): - pipeline_tasks.check_staleness("config.yaml") - mock_repo_cls.return_value.close.assert_called_once() + def test_dags_have_no_all_done_tasks(self) -> None: + # Airflow takes a run's state from its leaf tasks. A leaf with + # trigger_rule="all_done" succeeds after the pipeline task fails, which + # marks the whole run green and hides the failure. + dags = Path(__file__).resolve().parents[1] / "dags" + for dag_file in sorted(dags.glob("*_dag.py")): + self.assertNotIn("all_done", dag_file.read_text(encoding="utf-8"), dag_file.name) if __name__ == "__main__":