From 583d832511de32800d415d92acb112d148c9ce2c Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Thu, 5 Mar 2026 23:06:11 +0300 Subject: [PATCH] =?UTF-8?q?docs(airflow):=20=D0=BE=D0=B1=D0=BD=D0=BE=D0=B2?= =?UTF-8?q?=D0=BB=D0=B5=D0=BD=D0=BE=20=D1=80=D1=83=D0=BA=D0=BE=D0=B2=D0=BE?= =?UTF-8?q?=D0=B4=D1=81=D1=82=D0=B2=D0=BE=20=D0=BF=D0=BE=20=D0=BF=D1=80?= =?UTF-8?q?=D0=BE=D0=B3=D1=80=D0=B0=D0=BC=D0=BC=D0=BD=D0=BE=D0=BC=D1=83=20?= =?UTF-8?q?=D1=82=D0=B5=D1=81=D1=82=D0=B8=D1=80=D0=BE=D0=B2=D0=B0=D0=BD?= =?UTF-8?q?=D0=B8=D1=8E=20DAG?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - старая инструкция была неполной (без ODS/DDS/DM) и содержала избыточные требования. - Что: - переписан docs/agent-dag-testing.md с фокусом на итеративную разработку и Airflow REST API. - добавлены шаги по отладке упавших задач (логгирование через CLI) и проверке DWH слоев. - Проверка: - визуальная проверка текста руководства на соответствие актуальному пайплайну. --- docs/agent-dag-testing.md | 376 ++++++-------------------------------- 1 file changed, 51 insertions(+), 325 deletions(-) diff --git a/docs/agent-dag-testing.md b/docs/agent-dag-testing.md index d973fa6..7cf6d9e 100644 --- a/docs/agent-dag-testing.md +++ b/docs/agent-dag-testing.md @@ -1,363 +1,89 @@ # Тестирование DAG (гайд для AI-агентов) -Как программно проверить, что DAG работает корректно. -Все проверки выполняются через CLI, `docker compose exec` и SQL-запросы. -Браузер и Airflow UI **не используются**. +Этот документ описывает, как агенту взаимодействовать с Airflow (без UI) для точечной проверки и отладки DAG'ов в процессе разработки. + +## 0. Состояние среды (Предварительная проверка) + +Прежде чем тестировать DAG, убедитесь, что стек работает: +```bash +docker compose ps +``` +Если контейнеров нет, поднимите стек: `make up`. +Если исходная база данных пуста (например, после `make clean`), проинициализируйте её: `make bookings-init`. --- -## Предварительные условия - -Перед тестированием DAG стек должен быть поднят и здоров. +## 1. Быстрая проверка структуры графа (Локально) +При любом изменении Python-кода DAG'а сначала проверьте, что он компилируется и структура графа корректна: ```bash -# 1. Поднять стек (если не поднят) -make up - -# 2. Дождаться healthy-статуса всех сервисов -docker compose ps # greenplum и airflow-webserver должны быть (healthy) - -# 3. Для DAG bookings_to_gp_stage — инициализировать данные -make bookings-init # создать демо-БД bookings в контейнере bookings-db -make ddl-gp # создать STG-слой и внешние PXF-таблицы в Greenplum +make test ``` - -Если стек ранее сносился (`make clean`), шаги 1-3 обязательны. - -Перед запуском `bookings_to_gp_stage` обязательно проверьте, что source непустой: - -```bash -# Все значения ниже должны быть > 0 -docker compose exec bookings-db \ - psql -U bookings -d demo -At -c "SELECT COUNT(*) FROM bookings.bookings;" - -docker compose exec bookings-db \ - psql -U bookings -d demo -At -c "SELECT COUNT(*) FROM bookings.airports_data;" - -docker compose exec bookings-db \ - psql -U bookings -d demo -At -c "SELECT COUNT(*) FROM bookings.airplanes_data;" -``` - -Если хотя бы один `COUNT(*) = 0`, **не запускайте DAG**: -1. Выполните `make bookings-init`. -2. Повторите проверки `COUNT(*)`. -3. Если `bookings.bookings` всё ещё пустая, выполните `make bookings-generate-day` и проверьте снова. +Это запустит smoke-тесты (`tests/test_dags_smoke.py`), которые проверят целостность всех DAG'ов без обращения к базе данных. --- -## Уровень 1. Локальные проверки (без Docker) +## 2. Запуск конкретного DAG'а (REST API) -Быстрые проверки, не требующие поднятого стека: +Для тестирования загрузки данных запустите измененный DAG через REST API (базовый URL `http://localhost:8080/api/v1`, креды взять из `.venv`). +**Шаг 2.1. Снять DAG с паузы (если он новый):** ```bash -make test # pytest: unit-тесты helpers + smoke-тесты структуры DAG -make lint # black + isort в режиме проверки +curl -s -X PATCH "http://localhost:8080/api/v1/dags/" \ + -u admin:admin -H "Content-Type: application/json" -d '{"is_paused": false}' ``` -Smoke-тесты DAG (`tests/test_dags_smoke.py`) проверяют: -- DAG импортируется без ошибок; -- все ожидаемые `task_id` присутствуют; -- прямые рёбра графа совпадают с эталонными; -- задачи достижимы друг из друга (транзитивно). +**Шаг 2.2. Запустить DAG:** +```bash +curl -s -X POST "http://localhost:8080/api/v1/dags//dagRuns" \ + -u admin:admin -H "Content-Type: application/json" -d '{}' | jq '{dag_run_id}' +``` -Если Airflow не установлен в venv, smoke-тесты автоматически пропускаются (`skip`). +**Шаг 2.3. Проверить статус выполнения:** +Используйте `dag_run_id` из предыдущего шага: +```bash +curl -s "http://localhost:8080/api/v1/dags//dagRuns/" \ + -u admin:admin | jq '{state}' +``` +Повторяйте запрос, пока `state` не станет `success` или `failed`. --- -## Уровень 2. Тестовый прогон DAG (без записи в мета-БД) +## 3. Отладка упавших задач -Команда `airflow dags test` выполняет DAG целиком в оффлайн-режиме -(результат не сохраняется в Airflow, не создаётся `dag_run`): +Если DAG перешел в статус `failed`, найдите упавшую задачу: +**Шаг 3.1. Получить статусы всех задач:** ```bash -docker compose exec airflow-webserver \ - airflow dags test csv_to_greenplum 2024-01-01 - -docker compose exec airflow-webserver \ - airflow dags test bookings_to_gp_stage 2024-01-01 +curl -s "http://localhost:8080/api/v1/dags//dagRuns//taskInstances" \ + -u admin:admin | jq '.task_instances[] | {task_id, state}' ``` -Вывод идёт прямо в stdout — можно парсить на наличие `ERROR` / `FAILED`. +**Шаг 3.2. Посмотреть логи упавшей задачи (через CLI Airflow):** +REST API отдает логи сложно, поэтому для логов проще использовать `docker compose exec`: +```bash +docker compose exec airflow-webserver airflow tasks logs +``` + +*(Совет: ищите в логах слова `ERROR`, `Exception` или вывод SQL-ошибок от PostgresOperator).* --- -## Уровень 3. Полноценный запуск DAG (с записью в мета-БД) +## 4. Проверка результата в DWH (Greenplum) -### Запуск +Успешное выполнение DAG'а (зеленый статус) не гарантирует, что данные загрузились правильно (например, если источник был пуст). Проверьте целевые таблицы напрямую: ```bash -docker compose exec airflow-webserver \ - airflow dags trigger bookings_to_gp_stage +docker compose exec greenplum bash -lc "su - gpadmin -c \"/usr/local/greenplum-db/bin/psql -t -A -d gp_dwh -c 'SELECT COUNT(*) FROM <схема>.<таблица>;'\"" ``` - -Команда возвращает `run_id`. Если нужно получить его программно: - -```bash -docker compose exec airflow-webserver \ - airflow dags list-runs -d bookings_to_gp_stage -o json -``` - -### Ожидание завершения - -DAG может работать 30-60 секунд. Опрашиваем статус задач: - -```bash -docker compose exec airflow-webserver \ - airflow tasks states-for-dag-run bookings_to_gp_stage -o json -``` - -Повторять до тех пор, пока все задачи не перейдут в терминальный статус -(`success`, `failed`, `upstream_failed`, `skipped`). - -### Проверка результатов - -```bash -# Список задач и их статусы (текстовый формат) -docker compose exec airflow-webserver \ - airflow tasks states-for-dag-run bookings_to_gp_stage - -# Логи конкретной задачи (при отладке) -docker compose exec airflow-webserver \ - airflow tasks logs bookings_to_gp_stage load_airports_to_stg -``` - -**Критерий успеха:** все 20 задач в статусе `success`. +Убедитесь, что таблица содержит ожидаемое количество строк. --- -## Проверка параллельности - -В DAG `bookings_to_gp_stage` задачи `load_airports_to_stg` и `load_airplanes_to_stg` -должны запускаться параллельно (обе зависят только от `check_tickets_dq`). - -### Способ 1. По временным меткам (после реального запуска) +## 5. Полный сквозной тест (E2E) +Если вы вносили масштабные изменения, затрагивающие несколько слоев DWH, или меняли DDL таблиц, рекомендуется прогнать полный конвейер (STG -> ODS -> DDS -> DM): ```bash -docker compose exec airflow-webserver \ - airflow tasks states-for-dag-run bookings_to_gp_stage -o json +make e2e-etl ``` - -Сравнить `start_date` задач `load_airports_to_stg` и `load_airplanes_to_stg`. -**Критерий:** разница < 1 секунды. - -### Способ 2. По структуре графа (без запуска DAG) - -```bash -docker compose exec airflow-webserver python3 -c " -from airflow.models import DagBag - -dag = DagBag('/opt/airflow/dags').get_dag('bookings_to_gp_stage') - -airports = dag.get_task('load_airports_to_stg') -airplanes = dag.get_task('load_airplanes_to_stg') - -# Параллельность: задачи не зависят друг от друга -a_up = {t.task_id for t in airports.upstream_list} -b_up = {t.task_id for t in airplanes.upstream_list} - -print('airports upstream:', a_up) -print('airplanes upstream:', b_up) - -# airports не должен быть в upstream airplanes и наоборот -assert 'load_airports_to_stg' not in b_up, 'airplanes зависит от airports!' -assert 'load_airplanes_to_stg' not in a_up, 'airports зависит от airplanes!' -print('OK: задачи независимы, могут идти параллельно') -" -``` - -### Ожидаемые зависимости (эталон) - -| Задача | Ждёт (upstream) | -|--------|-----------------| -| `load_airports_to_stg` | `check_tickets_dq` | -| `load_airplanes_to_stg` | `check_tickets_dq` | -| `load_routes_to_stg` | `check_airports_dq` + `check_airplanes_dq` | -| `load_seats_to_stg` | `check_airplanes_dq` | -| `finish_summary` | `check_boarding_passes_dq` + `check_seats_dq` | - ---- - -## Проверка данных в Greenplum - -После успешного прогона DAG можно проверить наличие данных напрямую в БД: - -```bash -# Количество строк в ключевых таблицах -docker compose exec greenplum bash -lc \ - "su - gpadmin -c \"/usr/local/greenplum-db/bin/psql -t -A -d gp_dwh -c 'SELECT COUNT(*) FROM stg.bookings;'\"" - -docker compose exec greenplum bash -lc \ - "su - gpadmin -c \"/usr/local/greenplum-db/bin/psql -t -A -d gp_dwh -c 'SELECT COUNT(*) FROM stg.tickets;'\"" - -docker compose exec greenplum bash -lc \ - "su - gpadmin -c \"/usr/local/greenplum-db/bin/psql -t -A -d gp_dwh -c 'SELECT COUNT(*) FROM stg.airports;'\"" -``` - -**Критерий:** все таблицы непустые (COUNT > 0). - ---- - -## Проверка Airflow Connections - -Перед запуском DAG полезно убедиться, что подключения настроены: - -```bash -docker compose exec airflow-webserver airflow connections get greenplum_conn -docker compose exec airflow-webserver airflow connections get bookings_db -``` - -Обе команды должны вернуть параметры подключения без ошибок. - ---- - -## REST API (альтернатива CLI) - -REST API удобнее CLI для агента в ряде случаев: не нужен `docker exec`, -возвращает чистый JSON, проще поллить статус в цикле. - -**База:** `http://localhost:8080/api/v2` (порт из `AIRFLOW_WEB_PORT`, default: 8080) -**Аутентификация:** HTTP Basic Auth — `AIRFLOW_USER`/`AIRFLOW_PASSWORD` из `.env` (default: `admin`/`admin`) - -### Список DAG - -```bash -curl -s -u admin:admin http://localhost:8080/api/v2/dags | jq '.dags[].dag_id' -``` - -### Запуск DAG - -```bash -curl -s -u admin:admin \ - -X POST http://localhost:8080/api/v2/dags/bookings_to_gp_stage/dagRuns \ - -H "Content-Type: application/json" \ - -d '{}' | jq '{dag_run_id, state}' -``` - -Вернёт `dag_run_id` — он нужен для всех последующих запросов. - -### Статус запуска DAG - -```bash -curl -s -u admin:admin \ - http://localhost:8080/api/v2/dags/bookings_to_gp_stage/dagRuns/ \ - | jq '{state, start_date, end_date}' -``` - -Значения `state`: `queued` → `running` → `success` / `failed`. - -### Статусы всех задач запуска - -```bash -curl -s -u admin:admin \ - "http://localhost:8080/api/v2/dags/bookings_to_gp_stage/dagRuns//taskInstances" \ - | jq '.task_instances[] | {task_id, state, start_date}' -``` - -### Детали конкретной задачи - -```bash -curl -s -u admin:admin \ - "http://localhost:8080/api/v2/dags/bookings_to_gp_stage/dagRuns//taskInstances/load_airports_to_stg" \ - | jq '{task_id, state, start_date, end_date, duration}' -``` - -### Последний `dag_run_id` без явного сохранения - -```bash -curl -s -u admin:admin \ - "http://localhost:8080/api/v2/dags/bookings_to_gp_stage/dagRuns?order_by=-start_date&limit=1" \ - | jq -r '.dag_runs[0].dag_run_id' -``` - -### Когда использовать REST API вместо CLI - -| Ситуация | Предпочтительный способ | -|----------|------------------------| -| Нужен чистый JSON для парсинга | REST API | -| Агент работает вне Docker-хоста | REST API | -| Поллинг статуса в цикле | REST API (проще, чем `exec`) | -| Быстрая отладка или разовая проверка | CLI (`airflow dags test`) | -| Тестовый прогон без записи в мета-БД | CLI (`airflow dags test`) | - ---- - -## Полный E2E-тест (автоматизированный) - -Скрипт `scripts/e2e_smoke.sh` выполняет полный цикл: - -1. `make clean` — полный reset стека; -2. `make up` — поднимает сервисы; -3. ждёт `airflow-webserver` и `airflow-scheduler`; -4. `make bookings-init` — инициализирует демо-БД; -5. `make ddl-gp` — применяет DDL; -6. `make test` — локальные тесты; -7. `airflow dags test csv_to_greenplum 2024-01-01` — тест CSV-пайплайна; -8. проверяет `public.orders` непустую; -9. `airflow dags test bookings_to_gp_stage 2024-01-01` — тест bookings-пайплайна; -10. проверяет `stg.bookings` непустую. - -Запуск: - -```bash -./scripts/e2e_smoke.sh -``` - ---- - -## Список DAG и ожидаемые задачи - -### `csv_to_greenplum` (4 задачи) - -`create_orders_table` → `generate_csv` → `preview_csv` → `load_csv_to_greenplum` - -### `csv_to_greenplum_dq` (5 задач) - -`check_orders_table_exists` → `check_orders_schema` → `check_orders_has_rows` -→ `check_order_duplicates` → `data_quality_summary` - -### `bookings_to_gp_stage` (20 задач) - -``` -generate_bookings_day → load_bookings → check_bookings_dq - → load_tickets → check_tickets_dq - ├─ load_airports → check_airports_dq ─┐ - │ ├─ load_routes → check_routes_dq - ├─ load_airplanes → check_airplanes_dq ┤ → load_flights → check_flights_dq - │ │ → load_segments → check_segments_dq - │ │ → load_boarding_passes → check_bp_dq ─┐ - │ └─ load_seats → check_seats_dq ──────────────────────┤ - │ ▼ - └──────────────────────────────────────────────────────────────── finish_summary -``` - ---- - -## Ключевые команды (шпаргалка) - -| Действие | Команда | -|----------|---------| -| Список DAG | `docker compose exec airflow-webserver airflow dags list` | -| Список задач DAG | `docker compose exec airflow-webserver airflow tasks list ` | -| Тестовый прогон | `docker compose exec airflow-webserver airflow dags test 2024-01-01` | -| Запуск DAG | `docker compose exec airflow-webserver airflow dags trigger ` | -| Список запусков | `docker compose exec airflow-webserver airflow dags list-runs -d -o json` | -| Статусы задач | `docker compose exec airflow-webserver airflow tasks states-for-dag-run ` | -| Логи задачи | `docker compose exec airflow-webserver airflow tasks logs ` | -| Проверка подключений | `docker compose exec airflow-webserver airflow connections get ` | -| Запрос в Greenplum | `docker compose exec greenplum bash -lc "su - gpadmin -c '/usr/local/greenplum-db/bin/psql -t -A -d gp_dwh -c \"\"'"` | -| Здоровье стека | `docker compose ps` | - ---- - -## Типичные проблемы - -| Симптом | Вероятная причина | Что делать | -|---------|-------------------|------------| -| DAG не найден в `dags list` | Синтаксическая ошибка в файле | Посмотреть `docker compose logs airflow-scheduler` | -| `upstream_failed` у задачи | Упала задача выше по графу | Найти первую `failed`-задачу и смотреть её логи | -| DQ-проверка падает | Нет данных в source или нарушена целостность | Проверить данные в `bookings-db` и `stg.*` | -| `Connection ... not found` | Не задана переменная `AIRFLOW_CONN_*` | Проверить `.env` и `docker-compose.yml` | -| Greenplum `unhealthy` | PXF не стартовал (долгая инициализация) | Подождать 2-3 минуты, проверить `docker compose ps` | -| `relation ... does not exist` | Не применён DDL | Выполнить `make ddl-gp` | -| Пустые таблицы stg | Не выполнен `make bookings-init` | Выполнить `make bookings-init`, затем перезапустить DAG | -| `check_airports_dq` / `check_airplanes_dq` падают с `..._ext нет строк` | Source-таблицы в `bookings-db` пустые | Проверить `COUNT(*)` в `bookings.bookings`, `bookings.airports_data`, `bookings.airplanes_data`; затем `make bookings-init`/`make bookings-generate-day` | +Этот скрипт сам запустит все нужные DAG'и в правильном порядке и проверит результаты.