- Зачем:
- список DAG'ов должен читаться лесенкой ddl_init → world_init →
world_next_day, а путь менти — проходиться пустыми формами
(issue #4, спека редизайна пути менти, решения 2–3).
- Что:
- generator_control переименован в world_init, дефолт операции —
import; next-day ушёл из выпадашки в отдельный DAG;
- новый беспараметрный world_next_day: расписание */30 * * * *,
создаётся на паузе, catchup=False, max_active_runs=1; общие
задачи вынесены в airflow/dags/utils/startup_history_tasks.py;
- доки и контрактные тесты обновлены синхронно; быстрый старт
README — без make ddl, схему создаёт DAG ddl_init.
- Проверка:
- make test (210 + 31) и make lint зелёные;
- живая приёмка на чистом стенде: world_init пустой формой
импортировал эталонный мир за 217 с (3 дня, 280 437 событий),
world_next_day после снятия с паузы добавляет ровно один день
за прогон, дашборд Superset собирается.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
273 lines
9.0 KiB
Python
273 lines
9.0 KiB
Python
"""
|
|
DAG-пульт генератора стартовой истории.
|
|
|
|
Пульт выполняет только Python-код генератора внутри Airflow worker. Жизненный
|
|
цикл контейнеров остаётся в Makefile.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from datetime import datetime, timedelta
|
|
|
|
from airflow import DAG
|
|
from airflow.models.param import Param
|
|
from airflow.operators.empty import EmptyOperator
|
|
from airflow.operators.python import BranchPythonOperator, PythonOperator
|
|
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
|
|
from airflow.utils.trigger_rule import TriggerRule
|
|
from airflow_clickhouse_plugin.hooks.clickhouse import ClickHouseHook
|
|
|
|
from clickstream_generator.airflow_control import (
|
|
KAFKA_BOOTSTRAP_SERVERS,
|
|
assert_live_generator_not_running,
|
|
assert_stand_clean,
|
|
build_control_env,
|
|
default_artifact_path,
|
|
run_backfill,
|
|
run_import,
|
|
validate_import_artifact,
|
|
)
|
|
from clickstream_generator.launch import PROFILES
|
|
from utils.startup_history_tasks import (
|
|
assert_target_dag_not_paused,
|
|
check_manifest_task,
|
|
)
|
|
|
|
|
|
default_args = {
|
|
"owner": "airflow",
|
|
"depends_on_past": False,
|
|
"email_on_failure": False,
|
|
"email_on_retry": False,
|
|
"retries": 0,
|
|
"retry_delay": timedelta(minutes=1),
|
|
}
|
|
|
|
ETL_DAG_ID = "etl_pipeline"
|
|
|
|
|
|
def _conf(context) -> dict:
|
|
dag_run = context.get("dag_run")
|
|
return dag_run.conf if dag_run and dag_run.conf else {}
|
|
|
|
|
|
def _param(context, name: str):
|
|
return _conf(context).get(name, context["params"][name])
|
|
|
|
|
|
def _operation(context) -> str:
|
|
return str(_param(context, "operation"))
|
|
|
|
|
|
def _artifact_path(context) -> str:
|
|
value = str(_param(context, "artifact_path") or "").strip()
|
|
if value:
|
|
return value
|
|
return default_artifact_path(_operation(context))
|
|
|
|
|
|
def _overrides(context) -> dict[str, str]:
|
|
return {
|
|
key: str(_param(context, param_name) or "").strip()
|
|
for param_name, key in {
|
|
"seed": "GEN_SEED",
|
|
"model_time_speed": "GEN_MODEL_TIME_SPEED",
|
|
}.items()
|
|
}
|
|
|
|
|
|
def _build_env(context, operation: str) -> dict[str, str]:
|
|
artifact_path = _artifact_path(context)
|
|
if (
|
|
operation == "backfill"
|
|
and not str(_param(context, "artifact_path") or "").strip()
|
|
):
|
|
artifact_path = None
|
|
return build_control_env(
|
|
operation,
|
|
profile_name=str(_param(context, "profile")),
|
|
duration=str(_param(context, "duration") or "").strip(),
|
|
artifact_path=artifact_path,
|
|
overrides=_overrides(context),
|
|
)
|
|
|
|
|
|
def choose_operation(**context) -> str:
|
|
"""Выбирает ветку пульта по параметру operation."""
|
|
operation = _operation(context)
|
|
if operation == "backfill":
|
|
return "check_etl_not_paused_before_backfill"
|
|
if operation == "import":
|
|
return "check_etl_not_paused_before_import"
|
|
if operation == "check":
|
|
return "check_only"
|
|
raise ValueError(f"Неизвестная операция: {operation}")
|
|
|
|
|
|
def precheck_backfill(**context) -> None:
|
|
"""Проверяет чистоту стенда перед backfill."""
|
|
assert_live_generator_not_running()
|
|
hook = ClickHouseHook(clickhouse_conn_id="clickhouse_default", database="default")
|
|
assert_stand_clean(KAFKA_BOOTSTRAP_SERVERS, hook)
|
|
context["ti"].xcom_push(
|
|
key="generator_env",
|
|
value=_build_env(context, "backfill"),
|
|
)
|
|
|
|
|
|
def run_backfill_task(**context) -> None:
|
|
"""Выполняет backfill через код генератора."""
|
|
env = context["ti"].xcom_pull(task_ids="precheck_backfill", key="generator_env")
|
|
run_backfill(env)
|
|
|
|
|
|
def precheck_import(**context) -> None:
|
|
"""Проверяет чистоту стенда и совместимость артефакта перед import."""
|
|
assert_live_generator_not_running()
|
|
hook = ClickHouseHook(clickhouse_conn_id="clickhouse_default", database="default")
|
|
assert_stand_clean(KAFKA_BOOTSTRAP_SERVERS, hook)
|
|
env = _build_env(context, "import")
|
|
artifact_path = _artifact_path(context)
|
|
validate_import_artifact(env, artifact_path)
|
|
context["ti"].xcom_push(key="generator_env", value=env)
|
|
context["ti"].xcom_push(key="artifact_path", value=artifact_path)
|
|
|
|
|
|
def run_import_task(**context) -> None:
|
|
"""Выполняет import портативного артефакта."""
|
|
ti = context["ti"]
|
|
env = ti.xcom_pull(task_ids="precheck_import", key="generator_env")
|
|
artifact_path = ti.xcom_pull(task_ids="precheck_import", key="artifact_path")
|
|
run_import(env, artifact_path)
|
|
|
|
|
|
# Context7, Airflow 2.10.5: Param поддерживает enum.
|
|
with DAG(
|
|
dag_id="world_init",
|
|
description="Пульт стартовой истории генератора",
|
|
default_args=default_args,
|
|
schedule=None,
|
|
start_date=datetime(2024, 1, 1),
|
|
catchup=False,
|
|
max_active_runs=1,
|
|
is_paused_upon_creation=True,
|
|
tags=["generator", "startup-history"],
|
|
params={
|
|
"operation": Param(
|
|
"import",
|
|
type="string",
|
|
enum=["backfill", "import", "check"],
|
|
title="Операция",
|
|
description=(
|
|
"Что сделать: импортировать артефакт, создать историю "
|
|
"или проверить витрины."
|
|
),
|
|
),
|
|
"profile": Param(
|
|
"daily-wave",
|
|
type="string",
|
|
enum=sorted(PROFILES),
|
|
title="Профиль",
|
|
description="Именованный набор настроек генератора.",
|
|
),
|
|
# "null" в type делает поля формы необязательными (Context7, Airflow 2.10.5).
|
|
"duration": Param(
|
|
None,
|
|
type=["null", "string"],
|
|
title="Длительность",
|
|
description="Например 6h или 2d. Пусто — взять длительность из профиля.",
|
|
),
|
|
"seed": Param(
|
|
None,
|
|
type=["null", "string"],
|
|
title="GEN_SEED",
|
|
description="Пусто — взять seed из профиля.",
|
|
),
|
|
"model_time_speed": Param(
|
|
None,
|
|
type=["null", "string"],
|
|
title="GEN_MODEL_TIME_SPEED",
|
|
description="Пусто — взять скорость модельного времени из профиля.",
|
|
),
|
|
"artifact_path": Param(
|
|
None,
|
|
type=["null", "string"],
|
|
title="Артефакт",
|
|
description=(
|
|
"Backfill: куда сохранить файл; пусто — не сохранять. "
|
|
"Import: что читать; пусто — эталонный мир из репозитория."
|
|
),
|
|
),
|
|
},
|
|
) as dag:
|
|
route = BranchPythonOperator(
|
|
task_id="choose_operation",
|
|
python_callable=choose_operation,
|
|
)
|
|
|
|
precheck_backfill_task = PythonOperator(
|
|
task_id="precheck_backfill",
|
|
python_callable=precheck_backfill,
|
|
)
|
|
backfill_task = PythonOperator(
|
|
task_id="run_backfill",
|
|
python_callable=run_backfill_task,
|
|
)
|
|
|
|
precheck_import_task = PythonOperator(
|
|
task_id="precheck_import",
|
|
python_callable=precheck_import,
|
|
)
|
|
import_task = PythonOperator(
|
|
task_id="run_import",
|
|
python_callable=run_import_task,
|
|
)
|
|
|
|
check_etl_not_paused_before_backfill = PythonOperator(
|
|
task_id="check_etl_not_paused_before_backfill",
|
|
python_callable=assert_target_dag_not_paused,
|
|
op_kwargs={"dag_id": ETL_DAG_ID},
|
|
)
|
|
check_etl_not_paused_before_import = PythonOperator(
|
|
task_id="check_etl_not_paused_before_import",
|
|
python_callable=assert_target_dag_not_paused,
|
|
op_kwargs={"dag_id": ETL_DAG_ID},
|
|
)
|
|
trigger_etl = TriggerDagRunOperator(
|
|
task_id="trigger_etl",
|
|
trigger_dag_id=ETL_DAG_ID,
|
|
conf={"full_refresh": True},
|
|
wait_for_completion=True,
|
|
allowed_states=["success"],
|
|
failed_states=["failed"],
|
|
poke_interval=30,
|
|
trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS,
|
|
)
|
|
|
|
check_after_etl = PythonOperator(
|
|
task_id="check_after_etl",
|
|
python_callable=check_manifest_task,
|
|
)
|
|
|
|
check_only = PythonOperator(
|
|
task_id="check_only",
|
|
python_callable=check_manifest_task,
|
|
)
|
|
|
|
done = EmptyOperator(
|
|
task_id="done",
|
|
trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS,
|
|
)
|
|
|
|
route >> [
|
|
check_etl_not_paused_before_backfill,
|
|
check_etl_not_paused_before_import,
|
|
check_only,
|
|
]
|
|
check_etl_not_paused_before_backfill >> precheck_backfill_task
|
|
check_etl_not_paused_before_import >> precheck_import_task
|
|
precheck_backfill_task >> backfill_task >> trigger_etl
|
|
precheck_import_task >> import_task >> trigger_etl
|
|
trigger_etl >> check_after_etl >> done
|
|
check_only >> done
|