12 KiB
12 KiB
План развития Airflow DAG'ов
Цель
Перевести оркестрацию ETL на Airflow так, чтобы пайплайн оставался устойчивым к "грязным" данным и соответствовал целям проекта: Kafka → ClickHouse (STG → ODS → DDS → DM) → витрины для BI.
Что важно учесть в текущем репозитории
- В
AGENTS.mdкак quick check ожидается DAGetl_pipeline. - В Airflow-контейнере сейчас нет Kafka CLI, поэтому
kafka-topics.shиkafka-console-producer.shизBashOperatorне используем. - DDL должен выполняться строго последовательно:
00 -> 10 -> 20 -> 30 -> 40. - Файл
sql/dds/30_ods_to_dds.sqlуже включает обе загрузки (dds.clickиdds.event), поэтому в MVP это одна task.
Архитектура оркестрации
DAG 1 (обязательный): ddl_init
schedule:None(только ручной запуск).catchup:False.max_active_runs:1.is_paused_upon_creation:True.tags:["ddl", "bootstrap", "clickhouse"].
DAG 2 (обязательный): kafka_load
schedule:None(ручной/экспериментальный запуск).catchup:False.max_active_runs:1.is_paused_upon_creation:True.tags:["kafka", "ingest", "experiments"].- Реализация: этап 2 (после MVP).
DAG 3 (обязательный): etl_pipeline
schedule:None(ручной запуск для демо).catchup:False.max_active_runs:1.tags:["etl", "clickhouse", "demo"].
DAG 4 (опциональный): dq_monitor
schedule:0 * * * *.catchup:False.tags:["dq", "monitoring"].
Дизайн DAG ddl_init
Params
verify_only: bool, defaultfalse(прогон только проверок без применения DDL).
Tasks
| Task ID | Что делает | Источник SQL/реализация |
|---|---|---|
check_clickhouse |
Проверка доступности CH (SELECT 1) |
PythonOperator + clickhouse-connect |
ddl_00_databases |
Создание БД | sql/ddl/00_databases.sql |
ddl_10_stg |
STG + Kafka Engine + MV | sql/ddl/stg/10_stg.sql |
ddl_20_ods |
ODS таблицы + drop legacy MV STG→ODS | sql/ddl/ods/20_ods.sql |
ddl_30_dds |
Таблицы DDS | sql/ddl/dds/30_dds.sql |
ddl_40_dm |
VIEW витрины DM | sql/ddl/dm/40_dm.sql |
verify_schema |
Проверка ключевых таблиц/VIEW | SQL-check |
Зависимости:
check_clickhouse >> ddl_00_databases >> ddl_10_stg >> ddl_20_ods >> ddl_30_dds >> ddl_40_dm >> verify_schema
Примечание:
ddl_initзапускается вручную: при первом bootstrap, послеdocker compose down -v, после изменений схемы.
Дизайн DAG kafka_load (отдельный независимый контур)
Params (через Trigger DAG with config)
limit: int, default50.full_load: bool, defaultfalse.reset_topics: bool, defaulttrue.load_browser: bool, defaulttrue.load_location: bool, defaulttrue.load_device: bool, defaulttrue.load_geo: bool, defaulttrue.
TaskGroup precheck
| Task ID | Что делает | Реализация |
|---|---|---|
check_kafka |
Проверка доступности Kafka broker | PythonOperator + kafka-python |
check_input_files |
Проверка наличия data/*_events.jsonl |
PythonOperator |
validate_load_params |
Валидация параметров загрузки (limit, full_load, флаги потоков) |
PythonOperator |
TaskGroup ingest
| Task ID | Что делает | Реализация |
|---|---|---|
prepare_topics |
reset/create топиков по reset_topics |
PythonOperator + AdminClient |
load_browser_events |
Публикация browser_events.jsonl |
PythonOperator + KafkaProducer |
load_location_events |
Публикация location_events.jsonl |
PythonOperator + KafkaProducer |
load_device_events |
Публикация device_events.jsonl |
PythonOperator + KafkaProducer |
load_geo_events |
Публикация geo_events.jsonl |
PythonOperator + KafkaProducer |
verify_publish_counts |
Проверка, что отправлено > 0 сообщений в выбранные потоки | PythonOperator (по XCom) |
Зависимости:
precheck >> prepare_topics >> [load_browser_events, load_location_events, load_device_events, load_geo_events] >> verify_publish_counts
Примечания:
- DAG намеренно независим от
etl_pipeline: можно запускать ingest отдельно для экспериментов. - Авто-триггер
etl_pipelineне включаем по умолчанию; при необходимости добавляется отдельным параметром позже.
Дизайн DAG etl_pipeline
Params (через Trigger DAG with config)
full_refresh: bool, defaulttrue.wait_stg_timeout_sec: int, default600.
TaskGroup precheck
| Task ID | Что делает | Реализация |
|---|---|---|
check_clickhouse |
Проверка доступности CH (SELECT 1) |
PythonOperator + clickhouse-connect |
check_schema_ready |
Проверка, что DDL уже применён (stg.browser_raw, ods.browser_event, dds.event, dm.v_events_enriched) |
SQL-check, fail fast |
TaskGroup transform
| Task ID | Что делает | Источник SQL |
|---|---|---|
wait_for_stg_data |
Ожидание строк в stg.*_raw |
SQL-check |
load_ods |
Batch STG → ODS (основные таблицы + *_errors) | sql/ods/20_stg_to_ods.sql |
check_ods_quality |
Базовые DQ-метрики ODS (ошибки/total) | SQL-check |
truncate_dds_click |
Очистка dds.click при full_refresh=true |
inline SQL |
truncate_dds_event |
Очистка dds.event при full_refresh=true |
inline SQL |
load_dds |
ODS → DDS | sql/dds/30_ods_to_dds.sql |
check_dds_integrity |
Проверка orphan событий | inline SQL |
load_dm_summary |
DDS → DM DQ summary | sql/dm/40_dds_to_dm.sql |
validate_dm_summary |
Проверка, что dm.dq_summary не пуста |
SQL-check |
Зависимости:
wait_for_stg_data >> load_ods >> check_ods_quality >> [truncate_dds_click, truncate_dds_event] >> load_dds >> check_dds_integrity >> load_dm_summary >> validate_dm_summary
Итоговая цепочка etl_pipeline
precheck >> transform
Взаимодействие DAG'ов
Базовый сценарий:
ddl_init -> kafka_load -> etl_pipeline
Экспериментальные сценарии:
kafka_loadотдельно: проверить разные наборы/параметры загрузки без запуска transform.etl_pipelineотдельно: повторно пересчитать DDS/DM по уже загруженным данным.
Техническая реализация (приземленно)
ClickHouse в Airflow
- Использовать
clickhouse-connectнапрямую в Python helper. - Брать параметры подключения из
conn_id = clickhouse_defaultчерезBaseHook.get_connection.
Kafka в Airflow
- Добавить
kafka-pythonвairflow/requirements.txt. - Использовать Python-код для:
- reset/create топиков;
- публикации строк из
.jsonl(1 строка = 1 message value).
Общие helper-функции
dags/utils/clickhouse_helpers.py:execute_sql(sql: str) -> Noneexecute_sql_file(path: str) -> Nonefetch_one(sql: str) -> tuple
dags/utils/kafka_helpers.py:prepare_topics(reset: bool) -> Noneload_jsonl(file_path: str, topic: str, limit: int, full_load: bool) -> intcheck_kafka_ready() -> None
Структура файлов
dags/
├── __init__.py
├── ddl_init_dag.py # отдельный DAG для DDL (обязателен)
├── kafka_load_dag.py # отдельный DAG для ingest в Kafka (обязателен)
├── etl_pipeline_dag.py # основной DAG STG -> ODS -> DDS -> DM (обязателен)
├── dq_monitor_dag.py # опциональный DAG мониторинга
└── utils/
├── __init__.py
├── clickhouse_helpers.py
└── kafka_helpers.py
Этапы внедрения
- Этап 1 (MVP, обязательно):
- Реализовать
ddl_initиetl_pipeline. - Для загрузки данных использовать существующий сценарий
make data. - Проверить путь
ddl_init -> make data -> etl_pipeline.
- Реализовать
- Этап 2:
- Реализовать отдельный DAG
kafka_loadнаkafka-python. - Перенести загрузку из
make dataвkafka_load(функциональный паритет). - Добавить в
kafka_loadрасширенные параметры экспериментов (выбор потоков, сценарии reset/no-reset). - Добавить опциональный параметр автотриггера
etl_pipeline(по умолчаниюfalse).
- Реализовать отдельный DAG
- Этап 3:
- Добавить
dq_monitorи alert callback (email/Slack/webhook).
- Добавить
Критерии готовности
- Этап 1:
- В Airflow UI видны DAG
ddl_initиetl_pipeline. etl_pipelineпадает с понятной ошибкой, если схема не применена или STG пуста.- После прогона
make data -> etl_pipeline:- в
ods.browser_eventесть строки; - в
dds.clickиdds.eventесть строки; dm.dq_summaryзаполнена.
- в
- В Airflow UI видны DAG
- Этап 2:
- В Airflow UI дополнительно виден DAG
kafka_load. kafka_loadсlimit=50завершает отправку сообщений без падений.- После прогона
kafka_load -> etl_pipelineрезультаты совпадают сmake data -> etl_pipeline.
- В Airflow UI дополнительно виден DAG
- Для всех этапов:
- При повторном запуске
etl_pipelineсfull_refresh=trueнет неконтролируемых дублей в DDS.
- При повторном запуске
Минимальные smoke-checks
# 1) Запуск инфраструктуры
make up
# 2) Открыть Airflow UI
# http://localhost:8080 (admin/admin)
# 3) Один раз запустить ddl_init
# Trigger DAG ddl_init (без config или {"verify_only": false})
# 4) Этап 1: загрузить данные текущим способом
make data
# 5) Запустить etl_pipeline
# {"full_refresh": true}
# 6) Проверка результатов
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM ods.browser_event"
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM dds.click"
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM dds.event"
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM dm.dq_summary"
# 7) Этап 2: после реализации kafka_load
# Trigger DAG kafka_load {"limit": 50, "full_load": false, "reset_topics": true}
# Trigger DAG etl_pipeline {"full_refresh": true}
Следующий шаг после MVP
- Разделить
sql/dds/30_ods_to_dds.sqlна два файла и распараллелитьload_dds_clickиload_dds_event. - Перейти с
full_refreshна watermark-инкремент.