Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 5 additions & 7 deletions services/data-ngin/dags/data_pipeline_dag.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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()
12 changes: 5 additions & 7 deletions services/data-ngin/dags/new_data_pipeline_dag.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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()
12 changes: 5 additions & 7 deletions services/data-ngin/dags/tiingo_data_dag.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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()
29 changes: 0 additions & 29 deletions services/data-ngin/src/data_ngin/application/pipeline_tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")

Expand Down Expand Up @@ -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)
14 changes: 8 additions & 6 deletions services/data-ngin/tests/test_prod_wiring.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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__":
Expand Down
Loading