diff --git a/airflow/dags/ddl_init_dag.py b/airflow/dags/ddl_init_dag.py index ff34b0a..f3767c1 100644 --- a/airflow/dags/ddl_init_dag.py +++ b/airflow/dags/ddl_init_dag.py @@ -9,7 +9,6 @@ DAG инициализации DDL в ClickHouse. from __future__ import annotations -import re from datetime import datetime, timedelta from pathlib import Path @@ -20,6 +19,8 @@ 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 +from utils.airflow_params import parse_bool_param +from utils.sql_helpers import load_sql_statements as load_sql_file_statements # ----------------------------------------------------------------------------- @@ -55,23 +56,7 @@ SQL_ROOT = resolve_sql_root() 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) + return load_sql_file_statements(SQL_ROOT, relative_path) # ----------------------------------------------------------------------------- @@ -96,7 +81,10 @@ 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"])) + verify_only = parse_bool_param( + conf.get("verify_only", context["params"]["verify_only"]), + "verify_only", + ) return "skip_ddl" if verify_only else "ddl_00_databases" diff --git a/airflow/dags/etl_pipeline_dag.py b/airflow/dags/etl_pipeline_dag.py index ce43d2b..923fb5a 100644 --- a/airflow/dags/etl_pipeline_dag.py +++ b/airflow/dags/etl_pipeline_dag.py @@ -1,6 +1,9 @@ """ DAG ETL-процесса STG -> ODS -> DDS -> DM для учебного проекта. +Поток задач: + precheck -> transform: wait -> ods -> dq -> branch -> dds -> integrity -> dm -> validate + Принципы реализации: - SQL выполняется явными task на ClickHouseOperator; - SQL-файлы вызываются по фиксированным путям; @@ -9,7 +12,6 @@ DAG ETL-процесса STG -> ODS -> DDS -> DM для учебного про from __future__ import annotations -import re import time from datetime import datetime, timedelta from pathlib import Path @@ -23,6 +25,8 @@ 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 +from utils.airflow_params import parse_bool_param +from utils.sql_helpers import load_sql_statements as load_sql_file_statements # ----------------------------------------------------------------------------- @@ -58,23 +62,7 @@ SQL_ROOT = resolve_sql_root() 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) + return load_sql_file_statements(SQL_ROOT, relative_path) # ----------------------------------------------------------------------------- @@ -228,7 +216,10 @@ 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"])) + full_refresh = parse_bool_param( + conf.get("full_refresh", context["params"]["full_refresh"]), + "full_refresh", + ) return "transform.truncate_dds_click" if full_refresh else "transform.skip_truncate" @@ -245,6 +236,22 @@ def assert_dm_summary_not_empty(**context) -> None: raise AirflowException("dm.dq_summary пуста после load_dm_summary.") +def assert_dds_integrity(**context) -> None: + """Падает, если события ссылаются на отсутствующие клики.""" + ti = context["ti"] + result = ti.xcom_pull(task_ids="transform.check_dds_integrity") + + if not result or not result[0] or len(result[0]) != 1: + raise AirflowException(f"Некорректный результат check_dds_integrity: {result}") + + orphan_events = int(result[0][0]) + if orphan_events > 0: + raise AirflowException( + f"DDS integrity check failed: orphan_events={orphan_events}. " + "Есть события, чей click_id отсутствует в dds.click." + ) + + with DAG( dag_id="etl_pipeline", description="ETL STG -> ODS -> DDS -> DM для demo-проекта", @@ -342,6 +349,12 @@ with DAG( database="default", ) + assert_dds_integrity_task = PythonOperator( + task_id="assert_dds_integrity", + python_callable=assert_dds_integrity, + retries=0, + ) + load_dm_summary = ClickHouseOperator( task_id="load_dm_summary", sql=load_sql_statements("dm/40_dds_to_dm.sql"), @@ -364,6 +377,14 @@ with DAG( wait_for_stg_data_task >> load_ods >> 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 + ( + truncate_complete + >> load_dds + >> check_dds_integrity + >> assert_dds_integrity_task + >> load_dm_summary + >> validate_dm_summary_sql + >> validate_dm_summary + ) precheck >> transform diff --git a/airflow/dags/kafka_load_dag.py b/airflow/dags/kafka_load_dag.py index 484f8c8..f3f2783 100644 --- a/airflow/dags/kafka_load_dag.py +++ b/airflow/dags/kafka_load_dag.py @@ -24,6 +24,7 @@ from airflow.operators.python import PythonOperator from airflow.utils.task_group import TaskGroup # Импортируем helper-функции +from utils.airflow_params import parse_bool_param from utils.kafka_helpers import ( check_input_files, check_kafka_ready, @@ -69,7 +70,10 @@ def _validate_params(**context) -> None: def _prepare_topics(**context) -> None: """Подготовка топиков Kafka (создание/сброс).""" conf = context.get("dag_run", {}).conf or {} - reset_topics = bool(conf.get("reset_topics", context["params"]["reset_topics"])) + reset_topics = parse_bool_param( + conf.get("reset_topics", context["params"]["reset_topics"]), + "reset_topics", + ) prepare_topics(reset=reset_topics) diff --git a/airflow/dags/utils/airflow_params.py b/airflow/dags/utils/airflow_params.py new file mode 100644 index 0000000..ce232bc --- /dev/null +++ b/airflow/dags/utils/airflow_params.py @@ -0,0 +1,33 @@ +"""Общие helper-функции для параметров Airflow DAG.""" + +from __future__ import annotations + +TRUE_VALUES = {"1", "true", "t", "yes", "y", "on"} +FALSE_VALUES = {"0", "false", "f", "no", "n", "off"} + + +def _param_error(message: str) -> Exception: + try: + from airflow.exceptions import AirflowException + + return AirflowException(message) + except ModuleNotFoundError: + return ValueError(message) + + +def parse_bool_param(value: object, name: str) -> bool: + """Преобразует bool-параметр из dag_run.conf/params в явный boolean.""" + if isinstance(value, bool): + return value + + if isinstance(value, int) and value in (0, 1): + return bool(value) + + if isinstance(value, str): + normalized = value.strip().lower() + if normalized in TRUE_VALUES: + return True + if normalized in FALSE_VALUES: + return False + + raise _param_error(f"Параметр {name} должен быть boolean, получено: {value!r}") diff --git a/airflow/dags/utils/sql_helpers.py b/airflow/dags/utils/sql_helpers.py new file mode 100644 index 0000000..1b2b42f --- /dev/null +++ b/airflow/dags/utils/sql_helpers.py @@ -0,0 +1,153 @@ +"""Общие helper-функции для чтения SQL-файлов из Airflow DAG.""" + +from __future__ import annotations + +from pathlib import Path + + +def _sql_error(message: str) -> Exception: + try: + from airflow.exceptions import AirflowException + + return AirflowException(message) + except ModuleNotFoundError: + return ValueError(message) + + +def strip_sql_comments(sql_text: str) -> str: + """Удаляет SQL-комментарии, не трогая строки в кавычках.""" + result: list[str] = [] + i = 0 + in_single_quote = False + in_double_quote = False + + while i < len(sql_text): + char = sql_text[i] + next_char = sql_text[i + 1] if i + 1 < len(sql_text) else "" + + if in_single_quote: + result.append(char) + if char == "'" and next_char == "'": + result.append(next_char) + i += 2 + continue + if char == "'" and (i == 0 or sql_text[i - 1] != "\\"): + in_single_quote = False + i += 1 + continue + + if in_double_quote: + result.append(char) + if char == '"' and (i == 0 or sql_text[i - 1] != "\\"): + in_double_quote = False + i += 1 + continue + + if char == "'": + in_single_quote = True + result.append(char) + i += 1 + continue + + if char == '"': + in_double_quote = True + result.append(char) + i += 1 + continue + + if char == "-" and next_char == "-": + i += 2 + while i < len(sql_text) and sql_text[i] not in "\r\n": + i += 1 + continue + + if char == "/" and next_char == "*": + i += 2 + while ( + i < len(sql_text) - 1 + and not (sql_text[i] == "*" and sql_text[i + 1] == "/") + ): + if sql_text[i] in "\r\n": + result.append(sql_text[i]) + i += 1 + i += 2 + continue + + result.append(char) + i += 1 + + return "".join(result) + + +def split_sql_statements(sql_text: str) -> tuple[str, ...]: + """Делит SQL на команды по ';' вне строковых литералов.""" + statements: list[str] = [] + current: list[str] = [] + cleaned_sql = strip_sql_comments(sql_text) + in_single_quote = False + in_double_quote = False + i = 0 + + while i < len(cleaned_sql): + char = cleaned_sql[i] + next_char = cleaned_sql[i + 1] if i + 1 < len(cleaned_sql) else "" + + if in_single_quote: + current.append(char) + if char == "'" and next_char == "'": + current.append(next_char) + i += 2 + continue + if char == "'" and (i == 0 or cleaned_sql[i - 1] != "\\"): + in_single_quote = False + i += 1 + continue + + if in_double_quote: + current.append(char) + if char == '"' and (i == 0 or cleaned_sql[i - 1] != "\\"): + in_double_quote = False + i += 1 + continue + + if char == "'": + in_single_quote = True + current.append(char) + i += 1 + continue + + if char == '"': + in_double_quote = True + current.append(char) + i += 1 + continue + + if char == ";": + statement = "".join(current).strip() + if statement: + statements.append(statement) + current = [] + i += 1 + continue + + current.append(char) + i += 1 + + statement = "".join(current).strip() + if statement: + statements.append(statement) + + return tuple(statements) + + +def load_sql_statements(sql_root: Path, relative_path: str) -> tuple[str, ...]: + """Читает SQL-файл и возвращает отдельные команды.""" + file_path = sql_root / relative_path + if not file_path.is_file(): + raise _sql_error(f"SQL-файл не найден: {file_path}") + + statements = split_sql_statements(file_path.read_text(encoding="utf-8")) + if not statements: + raise _sql_error(f"SQL-файл пустой: {file_path}") + + return statements diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index ffb992e..e605eb3 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -59,6 +59,11 @@ - Запуск: ручной (`Trigger DAG with config`) - Параметр: `full_refresh` (`bool`, default `true`) — очистить DDS перед загрузкой - Зависимость: требует наличия данных в STG (от `kafka_load` или `make data`) +- Гейт целостности DDS: `check_dds_integrity` считает события без клика, а + `assert_dds_integrity` роняет DAG при `orphan_events > 0`. Проверка идёт после + `load_dds` и до `load_dm_summary`, чтобы DM не собирался поверх нарушенной связи + `dds.event -> dds.click`. Для `assert_dds_integrity` задано `retries=0`: повтор не + чинит уже собранную сироту и только задерживает явный failed-статус. ## Рекомендуемый сценарий (фаза 2) diff --git a/docs/course/LEARNING_PLAN.md b/docs/course/LEARNING_PLAN.md index 3b79f1a..fd460a9 100644 --- a/docs/course/LEARNING_PLAN.md +++ b/docs/course/LEARNING_PLAN.md @@ -113,15 +113,15 @@ Kafka из роадмапа — оно даёт словарь терминов. - Сами витрины DM показываем в деле в уроках 5–6, отдельного разбора не делаем. **Урок 4 — Airflow DAG (`airflow/dags/etl_pipeline_dag.py`):** -- [ ] **Главная правка (код↔доки):** сделать `check_dds_integrity` честным гейтом — +- [x] **Главная правка (код↔доки):** сделать `check_dds_integrity` честным гейтом — добавить `assert_dds_integrity` (PythonOperator), роняющий DAG при `orphan_events > 0`. Сейчас «проверка» только считает сирот в xcom и пишет их в `dq_summary`, но DAG остаётся зелёным, хотя `LESSON_STANDARD` §3 и `PRD` §2 (цель 3) обещают остановку при нарушении целостности. После правки доки и код сходятся. -- [ ] **Управляемая правка урока 4** строится на этом гейте: менти намеренно ломает +- [x] **Управляемая правка урока 4** строится на этом гейте: менти намеренно ломает целостность (вставляет «осиротевшее» событие) и видит, как DAG краснеет на `assert_dds_integrity`. Опирается на понятие сирот из урока 3. -- [ ] Добавить ASCII-поток задач в docstring DAG +- [x] Добавить ASCII-поток задач в docstring DAG (`precheck → transform: wait → ods → dq → branch → dds → integrity → dm → validate`) — для теста одного прохода. - Замечание: `check_ods_quality` тоже measure-only (метрики в xcom, без гейта) — это diff --git a/docs/course/README.md b/docs/course/README.md index 8560972..9abd56d 100644 --- a/docs/course/README.md +++ b/docs/course/README.md @@ -18,7 +18,7 @@ | [`PRD.md`](./PRD.md) | Рамка: зачем курс, цели, аудитория, скоуп, критерии успеха | Чтобы понять «что и зачем». Замороженный документ | | [`LEARNING_PLAN.md`](./LEARNING_PLAN.md) | План обучения: карта уроков, маршрут, аудит эталонных путей | Чтобы понять «в каком порядке и из чего» | | [`LESSON_STANDARD.md`](./LESSON_STANDARD.md) | Стандарт уроков: шаблон урока, качество кода, самопроверка | Рабочий чеклист при написании каждого урока | -| [`lessons/`](./lessons/) | Сами уроки, по одному файлу (есть: уроки 0–3) | Прохождение курса менти | +| [`lessons/`](./lessons/) | Сами уроки, по одному файлу (есть: уроки 0–4) | Прохождение курса менти | ## Порядок чтения diff --git a/docs/course/lessons/04_airflow_orchestration.md b/docs/course/lessons/04_airflow_orchestration.md new file mode 100644 index 0000000..c328202 --- /dev/null +++ b/docs/course/lessons/04_airflow_orchestration.md @@ -0,0 +1,374 @@ +# Урок 4. Оркестрация в Airflow + +> Формат: **практика** — будешь сам запускать команды и менять код, не только читать. +> Пререквизит: пройден урок 3 (слой DDS собран, ты знаешь, что такое событие-сирота и как +> посчитать `orphan_events`). +> Эталонный путь: [`airflow/dags/etl_pipeline_dag.py`](../../../airflow/dags/etl_pipeline_dag.py). +> +> Поток данных одной строкой: +> `precheck → transform: wait → ods → dq → branch → dds → integrity → dm → validate` +> +> О чём урок простыми словами: собираем все шаги STG → ODS → DDS → DM в один управляемый +> сценарий Airflow. И превращаем счётчик сирот из прошлого урока в стоп-кран: если событие +> ссылается на несуществующий клик, пайплайн краснеет и не идёт дальше. + +--- + +## 1. Зачем и где в проде + +В прошлых уроках мы смотрели на слои по отдельности: STG принял поток, ODS разобрал JSON, DDS +собрал сущности. Но в проде эти шаги не запускают «по памяти» руками. Иначе легко ошибиться: +сначала собрать DDS до ODS, забыть проверку, не заметить пустую витрину или пропустить сироту. + +Поэтому появляется **оркестрация** — управление порядком работ. В нашем стенде этим занимается +Airflow. Главная единица Airflow — **DAG** (Directed Acyclic Graph, направленный ациклический +граф). Проще: это схема задач без петли назад. В ней видно: + +- какие шаги есть в пайплайне; +- какой шаг ждёт какой; +- где пайплайн должен остановиться, если данные плохие; +- какой именно шаг покраснел, когда что-то сломалось. + +Отдельный шаг в Airflow называется **task** («задача»). Например, `load_ods` — одна task: +выполнить SQL для слоя ODS. `check_dds_integrity` — другая task: посчитать события-сироты. +А весь `etl_pipeline` — DAG, который связывает эти задачи в правильном порядке. + +### Проверка и гейт — не одно и то же + +Важная мысль урока: не каждая проверка должна ронять пайплайн. + +Есть проверки, которые **измеряют**. Например, `check_ods_quality` считает, сколько строк с +`parse_errors` появилось в ODS. Это полезная метрика: её можно отправить в лог, XCom или мониторинг. +Но сам факт ошибки парсинга в нашем учебном стенде не блокирует весь прогон. Грязные записи уже +отложены в `ods.*_errors`, чистые продолжают ехать дальше. + +А есть проверки, которые **гейтят**. **Гейт** (gate, «ворота») — это проверка, через которую +данные должны пройти, иначе следующие шаги не запускаются. Сироты в DDS — как раз такой случай. +Если `dds.event.click_id` ссылается на клик, которого нет в `dds.click`, целостность связей +сломана. В уроке 3 мы это только считали. Теперь DAG сам скажет: `orphan_events > 0`, дальше +не идём. + +> **В проде иначе.** Правил-гейтов обычно больше: свежесть данных, объём относительно вчера, +> доля ошибок, обязательные справочники. Но принцип тот же: одни проверки только измеряют, другие +> останавливают выпуск данных дальше. + +--- + +## 2. Руки: запускаем DAG и смотрим зелёный прогон + +Подними стенд, создай схему и залей малый срез: + +```bash +make up +make ddl +LIMIT=50 make data +``` + +Открой Airflow: `http://localhost:8080` (логин `admin`, пароль `admin`). Найди DAG +`etl_pipeline` и запусти его через **Trigger DAG with config**: + +```json +{"full_refresh": true} +``` + +`full_refresh=true` значит: перед сборкой DDS Airflow очистит `dds.click` и `dds.event`, а потом +заново наполнит их из ODS. Для учебного стенда это удобный чистый прогон: результат повторяемый, +старые эксперименты не мешают. + +Когда DAG завершится, открой его граф. На чистом срезе все задачи должны быть зелёными. Найди +внутри группы `transform` две задачи подряд: + +- `check_dds_integrity` — SQL-задача, которая считает сирот; +- `assert_dds_integrity` — Python-задача, которая решает, можно ли идти дальше. + +На чистом срезе `assert_dds_integrity` зелёная: сирот нет, пайплайн прошёл в DM. + +Проверь то же число в ClickHouse play-консоли `http://localhost:9123/play`: + +```sql +SELECT check_date, layer, table_name, check_name, check_value +FROM dm.dq_summary +WHERE layer = 'dds' + AND table_name = 'event_without_click' + AND check_name = 'orphan_events'; +``` + +Ожидаем `check_value = 0`. Это тот же смысл, что в уроке 3, только теперь число появилось внутри +управляемого прогона Airflow. + +--- + +## 3. Загляни внутрь + +Открой [`airflow/dags/etl_pipeline_dag.py`](../../../airflow/dags/etl_pipeline_dag.py). Сначала +прочитай верхний docstring. Там есть короткая карта задач: + +```text +precheck -> transform: wait -> ods -> dq -> branch -> dds -> integrity -> dm -> validate +``` + +Это не SQL-слои, а именно задачи DAG: + +- **`precheck`** — проверить, что ClickHouse доступен и схема уже создана; +- **`wait`** — дождаться строк в STG, чтобы не пересчитывать пустоту; +- **`ods`** — выполнить батч STG → ODS; +- **`dq`** — измерить качество ODS; +- **`branch`** — выбрать, чистить ли DDS перед загрузкой; +- **`dds`** — собрать `dds.click` и `dds.event`; +- **`integrity`** — проверить целостность событий и кликов; +- **`dm`** — собрать сводку качества `dm.dq_summary`; +- **`validate`** — убедиться, что сводка не пустая. + +Дальше разберём пять мест в файле. + +### `ClickHouseOperator`: SQL как отдельная задача + +Большая часть DAG — это задачи на `ClickHouseOperator`. Они выполняют SQL в ClickHouse: + +```python +load_ods = ClickHouseOperator( + task_id="load_ods", + sql=load_sql_statements("ods/20_stg_to_ods.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", +) +``` + +Это читается почти как команда: задача `load_ods` берёт SQL-файл `ods/20_stg_to_ods.sql` и +выполняет его в ClickHouse. Так же устроены `load_dds`, `load_dm_summary` и SQL-проверки. + +### `PythonOperator`: управляющая логика на Python + +**`PythonOperator`** — оператор Airflow, который запускает обычную Python-функцию как task. +В этом DAG Python нужен не для трансформации данных, а для управляющих решений: + +- `wait_for_stg_data` ждёт, пока в STG появятся строки; +- `assert_schema_ready` падает, если DDL не применён; +- `assert_dds_integrity` падает, если есть сироты; +- `assert_dm_summary_not_empty` падает, если финальная сводка пустая. + +То есть данные мы меняем SQL-ем, а Python оставляем для «можно ли продолжать». + +### `XCom`: маленькая передача результата между task + +Airflow хранит маленькие результаты задач в **XCom** (cross-communication, «передача между +задачами»). Это не место для данных пайплайна: туда не кладут таблицы и большие JSON. Но туда +нормально положить маленький результат проверки. + +Так работает пара `check_dds_integrity` → `assert_dds_integrity`: + +```python +check_dds_integrity = ClickHouseOperator( + task_id="check_dds_integrity", + sql=SQL_CHECK_DDS_INTEGRITY, + ... +) + +assert_dds_integrity_task = PythonOperator( + task_id="assert_dds_integrity", + python_callable=assert_dds_integrity, + retries=0, +) +``` + +SQL-задача возвращает одну строку с одним числом — `orphan_events`. Python-задача достаёт это +число из XCom: + +```python +result = ti.xcom_pull(task_ids="transform.check_dds_integrity") +``` + +Если число больше нуля, функция бросает `AirflowException`. Для Airflow это значит: task упала. +А у следующих задач по умолчанию правило «запускайся только если upstream успешен», поэтому +`load_dm_summary` дальше не пойдёт. + +У этой задачи отдельно стоит `retries=0`: если целостность уже нарушена, повтор через две минуты +ничего не исправит. Для учебного гейта честнее сразу показать красную задачу. + +> **Что проверили по API.** Перед правкой мы сверили Airflow через MCP Context7: `PythonOperator` +> подходит для Python-проверки, результат upstream можно читать через `ti.xcom_pull`, зависимости +> задаются оператором `>>`, а `AirflowException` переводит task в ошибку. Для ветвления после +> `BranchPythonOperator` оставлен `TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS`, чтобы join не +> пропускался из-за skipped-ветки. + +### `BranchPythonOperator`: развилка full refresh + +**Branching** («ветвление») — это выбор одной из нескольких веток DAG. У нас развилка простая: +чистить DDS перед загрузкой или не чистить. + +```python +choose_refresh_mode = BranchPythonOperator( + task_id="choose_refresh_mode", + python_callable=choose_full_refresh, +) +``` + +Если `full_refresh=true`, функция выбирает `transform.truncate_dds_click`: DAG очищает +`dds.click`, потом `dds.event`, и только после этого грузит DDS заново. Если `full_refresh=false`, +Airflow идёт через `skip_truncate` и сохраняет текущие строки DDS. + +После развилки ветки снова сходятся в `truncate_complete`. У этой задачи стоит специальное +правило `NONE_FAILED_MIN_ONE_SUCCESS`: «ни одна выбранная ветка не упала, и хотя бы одна успешно +прошла». Без него Airflow мог бы считать skipped-ветку проблемой. + +### Цепочка `>>`: порядок выполнения + +В самом низу файла порядок задач задан стрелками `>>`: + +```python +truncate_complete >> load_dds >> check_dds_integrity >> assert_dds_integrity_task +``` + +Читается слева направо: + +1. дождаться завершения развилки; +2. собрать DDS; +3. посчитать сирот; +4. проверить число и, если надо, уронить DAG. + +И только после этого идут `load_dm_summary`, `validate_dm_summary_sql` и `validate_dm_summary`. +Так гейт стоит именно там, где нужен: после сборки DDS, но до выпуска DM. + +--- + +## 4. Управляемая правка: заведём сироту и уроним DAG + +Сейчас воспроизведём обещание из урока 3: заведём событие-сироту и увидим, как Airflow красит +конкретную задачу. + +### Шаг 1. Вставь сироту в DDS + +Открой ClickHouse play-консоль `http://localhost:9123/play` и выполни: + +```sql +-- Событие со ссылкой на несуществующий клик dddd...-dddd (такого в dds.click нет) +INSERT INTO dds.event (event_id, event_ts, event_type, click_id, browser_name, dds_update_ts, ods_parse_errors) +VALUES ( + 'aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa', + now64(6), 'pageview', + 'dddddddd-dddd-dddd-dddd-dddddddddddd', + 'DemoBrowser', now64(3), [] +); +``` + +Проверь, что сирота правда появилась: + +```sql +SELECT count() AS orphans +FROM dds.event +WHERE click_id IS NOT NULL + AND click_id NOT IN (SELECT click_id FROM dds.click); +``` + +Ожидаем `orphans = 1`. + +### Шаг 2. Запусти DAG без очистки DDS + +Теперь вернись в Airflow и запусти `etl_pipeline` через **Trigger DAG with config**: + +```json +{"full_refresh": false} +``` + +Здесь важен именно `false`. Если поставить `true`, DAG сначала очистит `dds.event`, и наша +ручная сирота исчезнет ещё до проверки. А с `full_refresh=false` мы специально сохраняем текущий +DDS и даём гейту поймать битую связь. + +В графе Airflow ожидаем такую картину: + +- `check_dds_integrity` зелёная — SQL успешно посчитал `orphan_events`; +- `assert_dds_integrity` красная — Python-проверка увидела `orphan_events = 1` и бросила ошибку; +- `load_dm_summary` и последующие задачи дальше не пошли. + +Открой лог `assert_dds_integrity`. В нём должно быть сообщение примерно такого смысла: + +```text +DDS integrity check failed: orphan_events=1. Есть события, чей click_id отсутствует в dds.click. +``` + +Вот теперь сирота не просто «видна запросом». Она стала настоящим гейтом пайплайна. + +### Верни как было + +Чтобы вернуть стенд в чистое состояние, запусти `etl_pipeline` ещё раз, но уже с очисткой DDS: + +```json +{"full_refresh": true} +``` + +После зелёного прогона проверь: + +```sql +SELECT count() AS orphans +FROM dds.event +WHERE click_id IS NOT NULL + AND click_id NOT IN (SELECT click_id FROM dds.click); +``` + +Снова должно быть `0`. Если стенд после экспериментов совсем запутался, полный сброс остаётся +тем же: + +```bash +make clean && make up && make ddl && LIMIT=50 make data +``` + +После этого запусти `etl_pipeline` с `{"full_refresh": true}`. + +--- + +## 5. Проверь себя + +| Действие | Где смотреть | Что ожидать | +|----------|--------------|-------------| +| `etl_pipeline` с `{"full_refresh": true}` | Airflow graph | все задачи зелёные | +| чистый прогон | `dm.dq_summary`, строка `orphan_events` | `0` | +| вставка события-сироты | прямой SQL-счётчик сирот | `0 → 1` | +| `etl_pipeline` с `{"full_refresh": false}` после вставки | task `transform.assert_dds_integrity` | task красная, DAG failed | +| откат через `{"full_refresh": true}` | прямой SQL-счётчик сирот | снова `0` | + +--- + +## 6. Что должно получиться + +После урока у тебя на руках — видимый результат: + +- скрин Airflow graph, где `transform.assert_dds_integrity` красная после вставки сироты; +- и рядом короткое объяснение своими словами: почему `check_dds_integrity` зелёная, а + `assert_dds_integrity` красная. + +Проверь себя на словах — примерно эти вопросы всплывут на еженедельном созвоне: + +- что такое DAG и task в Airflow; +- чем проверка-метрика отличается от проверки-гейта; +- почему `check_ods_quality` только измеряет, а `assert_dds_integrity` останавливает пайплайн; +- зачем в управляемой правке нужен `full_refresh=false`; +- почему гейт стоит после `load_dds`, но до `load_dm_summary`. + +Если ответ на последний вопрос получается мутным, вернись к цепочке внизу DAG: DM должна +собираться только из DDS, который уже прошёл проверку целостности. + +--- + +## Вся цепочка разом: STG → ODS → DDS → DM + +Теперь у нас есть не только отдельные слои, но и порядок их жизни: + +- **STG** принимает поток и хранит сырой JSON; +- **ODS** типизирует и разделяет чистое/битое; +- **DDS** собирает сущности и проверяет связи между ними; +- **DM** даёт готовые витрины и сводку качества; +- **Airflow** связывает всё это в DAG, где виден порядок, статус и место падения. + +Главное изменение этого урока — сироты перестали быть ручным запросом «когда-нибудь посмотреть». +Теперь это автоматический гейт: если связь `dds.event → dds.click` порвалась, DAG показывает +красную задачу и не выпускает следующие шаги. + +--- + +## Мост к уроку 5 + +Airflow хорошо показывает судьбу конкретного прогона: зелёный он или красный, на какой задаче +упал, что написано в логе. Но есть другой вопрос: как увидеть состояние всего стенда со стороны — +живы ли сервисы, есть ли лаг в Kafka, не пропали ли метрики качества? В уроке 5 перейдём к +мониторингу: Prometheus и Grafana покажут пайплайн не как один DAG-run, а как систему, за которой +можно наблюдать постоянно. diff --git a/sql/dm/40_dds_to_dm.sql b/sql/dm/40_dds_to_dm.sql index f8f22a2..8b70592 100644 --- a/sql/dm/40_dds_to_dm.sql +++ b/sql/dm/40_dds_to_dm.sql @@ -15,7 +15,7 @@ -- чтобы не накапливать дубликаты при повторных прогонах. -- -- Витрины DM сейчас — это VIEW (логика без копии данных). Если тяжёлая агрегация --- начнёт тормозить, её материализуют в таблицу; пример — в docs/ARCHITECTURE.md, +-- начнёт тормозить, её материализуют в таблицу — пример в docs/ARCHITECTURE.md, -- раздел «Материализация витрин». -- ============================================================================