""" 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