diff --git a/Dockerfile.airflow b/Dockerfile.airflow index 0643813..b857115 100644 --- a/Dockerfile.airflow +++ b/Dockerfile.airflow @@ -1,5 +1,5 @@ # Use the official Airflow image as base -FROM apache/airflow:2.9.3 +FROM apache/airflow:2.10.5 # Set environment variables ENV AIRFLOW_HOME=/opt/airflow diff --git a/airflow/requirements.txt b/airflow/requirements.txt index cdeea91..feae4fb 100644 --- a/airflow/requirements.txt +++ b/airflow/requirements.txt @@ -1,10 +1,7 @@ -# Airflow requirements для ClickHouse DWH проекта +# Airflow requirements для учебного ETL-проекта -# Core database connector для metadata +# Metadata DB для Airflow psycopg2-binary==2.9.9 -# ClickHouse provider для ETL +# ClickHouse operator/hook для DAG'ов airflow-clickhouse-plugin==1.6.0 - -# Для работы с данными -pandas==2.1.4 diff --git a/dags/__pycache__/ddl_init_dag.cpython-312.pyc b/dags/__pycache__/ddl_init_dag.cpython-312.pyc new file mode 100644 index 0000000..9f0a801 Binary files /dev/null and b/dags/__pycache__/ddl_init_dag.cpython-312.pyc differ diff --git a/dags/__pycache__/etl_pipeline_dag.cpython-312.pyc b/dags/__pycache__/etl_pipeline_dag.cpython-312.pyc new file mode 100644 index 0000000..aa2e353 Binary files /dev/null and b/dags/__pycache__/etl_pipeline_dag.cpython-312.pyc differ diff --git a/dags/ddl_init_dag.py b/dags/ddl_init_dag.py new file mode 100644 index 0000000..9529177 --- /dev/null +++ b/dags/ddl_init_dag.py @@ -0,0 +1,188 @@ +""" +DAG инициализации DDL в ClickHouse. + +Учебный формат: +- каждая операция DDL выполняется отдельной SQL-task; +- SQL-файлы вызываются явно по фиксированным путям; +- режим verify_only позволяет прогонять только проверки схемы. +""" + +from __future__ import annotations + +import re +from datetime import datetime, timedelta +from pathlib import Path + +from airflow import DAG +from airflow.exceptions import AirflowException +from airflow.models.param import Param +from airflow.operators.empty import EmptyOperator +from airflow.operators.python import BranchPythonOperator, PythonOperator +from airflow.utils.trigger_rule import TriggerRule +from airflow_clickhouse_plugin.operators.clickhouse import ClickHouseOperator + + +# ----------------------------------------------------------------------------- +# Базовые настройки DAG +# ----------------------------------------------------------------------------- +default_args = { + "owner": "airflow", + "depends_on_past": False, + "email_on_failure": False, + "email_on_retry": False, + "retries": 1, + "retry_delay": timedelta(minutes=2), +} + + +# ----------------------------------------------------------------------------- +# SQL-файлы проекта +# ----------------------------------------------------------------------------- +SQL_ROOT = Path(__file__).resolve().parents[1] / "sql" + + +def load_sql_statements(relative_path: str) -> tuple[str, ...]: + """Читает SQL-файл и делит его на отдельные команды по ';'.""" + file_path = SQL_ROOT / relative_path + if not file_path.is_file(): + raise AirflowException(f"SQL-файл не найден: {file_path}") + + sql_text = file_path.read_text(encoding="utf-8") + statements: list[str] = [] + for segment in sql_text.split(";"): + # Убираем блочные и строковые комментарии, чтобы не отправлять "пустые" запросы. + no_block_comments = re.sub(r"/\*.*?\*/", "", segment, flags=re.S) + lines = [line for line in no_block_comments.splitlines() if not line.strip().startswith("--")] + cleaned = "\n".join(lines).strip() + if cleaned: + statements.append(cleaned) + + if not statements: + raise AirflowException(f"SQL-файл пустой: {file_path}") + return tuple(statements) + + +# ----------------------------------------------------------------------------- +# SQL-проверки +# ----------------------------------------------------------------------------- +SQL_CHECK_CLICKHOUSE = "SELECT 1 AS ok" + +SQL_VERIFY_SCHEMA = """ +SELECT + (SELECT count() FROM system.tables WHERE database = 'stg' AND name = 'browser_raw') AS stg_browser_raw, + (SELECT count() FROM system.tables WHERE database = 'ods' AND name = 'browser_event') AS ods_browser_event, + (SELECT count() FROM system.tables WHERE database = 'dds' AND name = 'click') AS dds_click, + (SELECT count() FROM system.tables WHERE database = 'dds' AND name = 'event') AS dds_event, + (SELECT count() FROM system.tables WHERE database = 'dm' AND name = 'v_events_enriched') AS dm_v_events_enriched +""" + + +# ----------------------------------------------------------------------------- +# Управляющие функции +# ----------------------------------------------------------------------------- +def choose_ddl_mode(**context) -> str: + """Выбирает ветку выполнения: full DDL или только verify.""" + dag_run = context.get("dag_run") + conf = dag_run.conf if dag_run else {} + verify_only = bool(conf.get("verify_only", context["params"]["verify_only"])) + return "skip_ddl" if verify_only else "ddl_00_databases" + + +def assert_schema_ready(**context) -> None: + """Проверяет результат финальной SQL-проверки схемы.""" + ti = context["ti"] + result = ti.xcom_pull(task_ids="verify_schema_sql") + + if not result or not result[0] or len(result[0]) != 5: + raise AirflowException(f"Некорректный результат проверки схемы: {result}") + + if any(value == 0 for value in result[0]): + raise AirflowException( + "Схема применена не полностью. Проверьте таблицы/VIEW stg, ods, dds, dm." + ) + + +with DAG( + dag_id="ddl_init", + description="Инициализация схемы ClickHouse (stg/ods/dds/dm)", + default_args=default_args, + schedule=None, + start_date=datetime(2024, 1, 1), + catchup=False, + max_active_runs=1, + is_paused_upon_creation=True, + tags=["ddl", "bootstrap", "clickhouse"], + params={ + "verify_only": Param(False, type="boolean"), + }, +) as dag: + check_clickhouse = ClickHouseOperator( + task_id="check_clickhouse", + sql=SQL_CHECK_CLICKHOUSE, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + choose_mode = BranchPythonOperator( + task_id="choose_mode", + python_callable=choose_ddl_mode, + ) + + ddl_00_databases = ClickHouseOperator( + task_id="ddl_00_databases", + sql=load_sql_statements("ddl/00_databases.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + ddl_10_stg = ClickHouseOperator( + task_id="ddl_10_stg", + sql=load_sql_statements("ddl/stg/10_stg.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + ddl_20_ods = ClickHouseOperator( + task_id="ddl_20_ods", + sql=load_sql_statements("ddl/ods/20_ods.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + ddl_30_dds = ClickHouseOperator( + task_id="ddl_30_dds", + sql=load_sql_statements("ddl/dds/30_dds.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + ddl_40_dm = ClickHouseOperator( + task_id="ddl_40_dm", + sql=load_sql_statements("ddl/dm/40_dm.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + skip_ddl = EmptyOperator(task_id="skip_ddl") + + ddl_complete = EmptyOperator( + task_id="ddl_complete", + trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS, + ) + + verify_schema_sql = ClickHouseOperator( + task_id="verify_schema_sql", + sql=SQL_VERIFY_SCHEMA, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + verify_schema = PythonOperator( + task_id="verify_schema", + python_callable=assert_schema_ready, + ) + + check_clickhouse >> choose_mode + choose_mode >> skip_ddl >> ddl_complete + choose_mode >> ddl_00_databases >> ddl_10_stg >> ddl_20_ods >> ddl_30_dds >> ddl_40_dm >> ddl_complete + ddl_complete >> verify_schema_sql >> verify_schema diff --git a/dags/etl_pipeline_dag.py b/dags/etl_pipeline_dag.py index feb0b1c..586f0dc 100644 --- a/dags/etl_pipeline_dag.py +++ b/dags/etl_pipeline_dag.py @@ -1,46 +1,282 @@ """ -ETL Pipeline DAG для ClickHouse DWH +DAG ETL-процесса ODS -> DDS -> DM для учебного проекта. -Шаблон DAG для оркестрации пайплайна данных. -Полная реализация будет добавлена позже. - -Пайплайн: - 1. DDL - создание структуры БД - 2. Load - загрузка данных в Kafka - 3. Transform - batch трансформация ODS → DDS → DM +Принципы реализации: +- SQL выполняется явными task на ClickHouseOperator; +- SQL-файлы вызываются по фиксированным путям; +- Python используется только для управляющей логики (branch/wait/assert). """ +from __future__ import annotations + +import re +import time from datetime import datetime, timedelta -from airflow import DAG -from airflow.operators.bash import BashOperator -from airflow.operators.empty import EmptyOperator +from pathlib import Path +from airflow import DAG +from airflow.exceptions import AirflowException +from airflow.models.param import Param +from airflow.operators.empty import EmptyOperator +from airflow.operators.python import BranchPythonOperator, PythonOperator +from airflow.utils.task_group import TaskGroup +from airflow.utils.trigger_rule import TriggerRule +from airflow_clickhouse_plugin.hooks.clickhouse import ClickHouseHook +from airflow_clickhouse_plugin.operators.clickhouse import ClickHouseOperator + + +# ----------------------------------------------------------------------------- # Базовые настройки DAG +# ----------------------------------------------------------------------------- default_args = { "owner": "airflow", "depends_on_past": False, "email_on_failure": False, "email_on_retry": False, "retries": 1, - "retry_delay": timedelta(minutes=5), + "retry_delay": timedelta(minutes=2), } + +# ----------------------------------------------------------------------------- +# SQL-файлы проекта +# ----------------------------------------------------------------------------- +SQL_ROOT = Path(__file__).resolve().parents[1] / "sql" + + +def load_sql_statements(relative_path: str) -> tuple[str, ...]: + """Читает SQL-файл и делит его на отдельные команды по ';'.""" + file_path = SQL_ROOT / relative_path + if not file_path.is_file(): + raise AirflowException(f"SQL-файл не найден: {file_path}") + + sql_text = file_path.read_text(encoding="utf-8") + statements: list[str] = [] + for segment in sql_text.split(";"): + # Убираем блочные и строковые комментарии, чтобы не отправлять "пустые" запросы. + no_block_comments = re.sub(r"/\*.*?\*/", "", segment, flags=re.S) + lines = [line for line in no_block_comments.splitlines() if not line.strip().startswith("--")] + cleaned = "\n".join(lines).strip() + if cleaned: + statements.append(cleaned) + + if not statements: + raise AirflowException(f"SQL-файл пустой: {file_path}") + return tuple(statements) + + +# ----------------------------------------------------------------------------- +# SQL для проверок и технических шагов +# ----------------------------------------------------------------------------- +SQL_CHECK_CLICKHOUSE = "SELECT 1 AS ok" + +SQL_CHECK_SCHEMA_READY = """ +SELECT + (SELECT count() FROM system.tables WHERE database = 'stg' AND name = 'browser_raw') AS stg_browser_raw, + (SELECT count() FROM system.tables WHERE database = 'ods' AND name = 'browser_event') AS ods_browser_event, + (SELECT count() FROM system.tables WHERE database = 'dds' AND name = 'event') AS dds_event, + (SELECT count() FROM system.tables WHERE database = 'dm' AND name = 'v_events_enriched') AS dm_v_events_enriched +""" + +SQL_CHECK_ODS_QUALITY = """ +SELECT + count() AS total_rows, + countIf(length(parse_errors) > 0) AS rows_with_errors, + round(if(count() = 0, 0, countIf(length(parse_errors) > 0) / count() * 100), 2) AS error_pct +FROM ods.browser_event +""" + +SQL_TRUNCATE_DDS_CLICK = "TRUNCATE TABLE dds.click" +SQL_TRUNCATE_DDS_EVENT = "TRUNCATE TABLE dds.event" + +SQL_CHECK_DDS_INTEGRITY = """ +SELECT + countIf(click_id IS NOT NULL AND click_id NOT IN (SELECT click_id FROM dds.click)) AS orphan_events +FROM dds.event +""" + +SQL_VALIDATE_DM_SUMMARY = "SELECT count() AS dq_rows FROM dm.dq_summary" + + +# ----------------------------------------------------------------------------- +# Управляющие функции +# ----------------------------------------------------------------------------- +def assert_schema_ready(**context) -> None: + """Падает, если DDL не применён полностью.""" + ti = context["ti"] + result = ti.xcom_pull(task_ids="precheck.check_schema_ready_sql") + + if not result or not result[0] or len(result[0]) != 4: + raise AirflowException(f"Некорректный результат check_schema_ready_sql: {result}") + + if any(value == 0 for value in result[0]): + raise AirflowException( + "Схема не готова: сначала запустите DAG ddl_init, затем повторите etl_pipeline." + ) + + +def wait_for_ods_data(**context) -> None: + """ + Ожидает появления строк в ods.browser_event до заданного таймаута. + Таймаут берётся из dag_run.conf.wait_ods_timeout_sec или из params. + """ + dag_run = context.get("dag_run") + conf = dag_run.conf if dag_run else {} + timeout_sec = int(conf.get("wait_ods_timeout_sec", context["params"]["wait_ods_timeout_sec"])) + poll_interval_sec = 10 + + hook = ClickHouseHook(clickhouse_conn_id="clickhouse_default", database="default") + started = time.monotonic() + + while True: + rows = hook.execute("SELECT count() FROM ods.browser_event") + count_rows = int(rows[0][0]) if rows else 0 + if count_rows > 0: + return + + elapsed = int(time.monotonic() - started) + if elapsed >= timeout_sec: + raise AirflowException( + f"Таймаут ожидания ODS истёк ({timeout_sec} сек). " + "Таблица ods.browser_event всё ещё пуста." + ) + + time.sleep(poll_interval_sec) + + +def choose_full_refresh(**context) -> str: + """Ветвление: делать TRUNCATE DDS или пропустить.""" + dag_run = context.get("dag_run") + conf = dag_run.conf if dag_run else {} + full_refresh = bool(conf.get("full_refresh", context["params"]["full_refresh"])) + return "transform.truncate_dds_click" if full_refresh else "transform.skip_truncate" + + +def assert_dm_summary_not_empty(**context) -> None: + """Проверяет, что dm.dq_summary заполнена после загрузки.""" + ti = context["ti"] + result = ti.xcom_pull(task_ids="transform.validate_dm_summary_sql") + + if not result or not result[0] or len(result[0]) != 1: + raise AirflowException(f"Некорректный результат validate_dm_summary_sql: {result}") + + dq_rows = int(result[0][0]) + if dq_rows <= 0: + raise AirflowException("dm.dq_summary пуста после load_dm_summary.") + + with DAG( dag_id="etl_pipeline", + description="ETL ODS -> DDS -> DM для demo-проекта", default_args=default_args, - description="ETL pipeline для ClickHouse DWH", - schedule=None, # Запуск только вручную (пока) + schedule=None, start_date=datetime(2024, 1, 1), catchup=False, - tags=["etl", "clickhouse", "dwh"], + max_active_runs=1, + is_paused_upon_creation=True, + tags=["etl", "clickhouse", "demo"], + params={ + "full_refresh": Param(True, type="boolean"), + "wait_ods_timeout_sec": Param(600, type="integer", minimum=30), + }, ) as dag: - - # TODO: добавить задачи пайплайна - # - ddl: создание структуры БД - # - load: загрузка данных в Kafka - # - transform: batch трансформация - - start = EmptyOperator(task_id="start") - end = EmptyOperator(task_id="end") - - start >> end + with TaskGroup(group_id="precheck") as precheck: + check_clickhouse = ClickHouseOperator( + task_id="check_clickhouse", + sql=SQL_CHECK_CLICKHOUSE, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + check_schema_ready_sql = ClickHouseOperator( + task_id="check_schema_ready_sql", + sql=SQL_CHECK_SCHEMA_READY, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + check_schema_ready = PythonOperator( + task_id="check_schema_ready", + python_callable=assert_schema_ready, + ) + + check_clickhouse >> check_schema_ready_sql >> check_schema_ready + + with TaskGroup(group_id="transform") as transform: + wait_for_ods_data_task = PythonOperator( + task_id="wait_for_ods_data", + python_callable=wait_for_ods_data, + ) + + check_ods_quality = ClickHouseOperator( + task_id="check_ods_quality", + sql=SQL_CHECK_ODS_QUALITY, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + choose_refresh_mode = BranchPythonOperator( + task_id="choose_refresh_mode", + python_callable=choose_full_refresh, + ) + + truncate_dds_click = ClickHouseOperator( + task_id="truncate_dds_click", + sql=SQL_TRUNCATE_DDS_CLICK, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + truncate_dds_event = ClickHouseOperator( + task_id="truncate_dds_event", + sql=SQL_TRUNCATE_DDS_EVENT, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + skip_truncate = EmptyOperator(task_id="skip_truncate") + + truncate_complete = EmptyOperator( + task_id="truncate_complete", + trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS, + ) + + load_dds = ClickHouseOperator( + task_id="load_dds", + sql=load_sql_statements("dds/30_ods_to_dds.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + check_dds_integrity = ClickHouseOperator( + task_id="check_dds_integrity", + sql=SQL_CHECK_DDS_INTEGRITY, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + load_dm_summary = ClickHouseOperator( + task_id="load_dm_summary", + sql=load_sql_statements("dm/40_dds_to_dm.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + validate_dm_summary_sql = ClickHouseOperator( + task_id="validate_dm_summary_sql", + sql=SQL_VALIDATE_DM_SUMMARY, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + validate_dm_summary = PythonOperator( + task_id="validate_dm_summary", + python_callable=assert_dm_summary_not_empty, + ) + + wait_for_ods_data_task >> check_ods_quality >> choose_refresh_mode + choose_refresh_mode >> truncate_dds_click >> truncate_dds_event >> truncate_complete + choose_refresh_mode >> skip_truncate >> truncate_complete + truncate_complete >> load_dds >> check_dds_integrity >> load_dm_summary >> validate_dm_summary_sql >> validate_dm_summary + + precheck >> transform diff --git a/docker-compose.yml b/docker-compose.yml index 5da65f9..976d03f 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -4,7 +4,7 @@ x-airflow-env: &airflow-default-env AIRFLOW__CORE__EXECUTOR: LocalExecutor AIRFLOW__WEBSERVER__SECRET_KEY: ${AIRFLOW_SECRET_KEY:-replace-me-with-random-string} # ClickHouse connection для ETL - AIRFLOW_CONN_CLICKHOUSE_DEFAULT: clickhouse://default:123456@clickhouse:8123/default + AIRFLOW_CONN_CLICKHOUSE_DEFAULT: clickhouse://default:123456@clickhouse:9000/default services: @@ -97,7 +97,7 @@ services: build: context: . dockerfile: Dockerfile.airflow - image: airflow-optimized:2.9.3 + image: airflow-optimized:2.10.5 environment: <<: *airflow-default-env command: > @@ -108,6 +108,7 @@ services: - "8080:8080" volumes: - ./dags:/opt/airflow/dags + - ./sql:/opt/airflow/sql:ro - ./data:/opt/airflow/data networks: - cs_dwh @@ -132,6 +133,7 @@ services: " volumes: - ./dags:/opt/airflow/dags + - ./sql:/opt/airflow/sql:ro - ./data:/opt/airflow/data networks: - cs_dwh @@ -153,13 +155,14 @@ services: <<: *airflow-default-env volumes: - ./dags:/opt/airflow/dags + - ./sql:/opt/airflow/sql:ro - ./data:/opt/airflow/data networks: - cs_dwh command: > bash -ceuo pipefail " mkdir -p /opt/airflow/data && - chmod -R 777 /opt/airflow/data || true && + chmod -R a+rX /opt/airflow/data || true && chown -R airflow:0 /opt/airflow/data || true && umask 000 && su -s /bin/bash airflow -c 'airflow db migrate' &&