6.0 KiB
6.0 KiB
Мини‑README по учебному DAG bookings_to_gp_stage (черновик)
Внутренний файл, чтобы не забыть договорённости. Перед итоговой сдачей документацию по блоку bookings/STG нужно будет аккуратно собрать и переписать.
1. Что делает DAG
- DAG
bookings_to_gp_stageпоказывает учебный поток:- источник: демо‑БД
bookings-db(Postgres, схемаbookings, таблицаbookings.bookings); - при каждом запуске генерируется один учебный день данных (идемпотентно);
- данные из
bookings.bookingsпереливаются в сырой слойstg.bookingsв Greenplum через PXF‑внешнюю таблицуstg.bookings_ext.
- источник: демо‑БД
- Слой
stgзадуман как «сырой»:- все бизнес‑колонки (
book_ref,book_date,total_amount) хранятся какTEXT; - есть тех.колонки
src_created_at_ts,load_dttm,batch_id.
- все бизнес‑колонки (
Подробный дизайн описан в docs/internal/bookings_stg_design.md.
2. Что нужно, чтобы DAG завёлся
Минимальные предпосылки:
- Стенд поднят:
make up. - Демо‑БД bookings инициализирована:
make bookings-init. - В Greenplum применён DDL (созданы схема
stgи таблицыstg.bookings_ext/stg.bookings):- учебный вариант: запустить DAG
bookings_stg_ddl(он используетsql/stg/bookings_ddl.sql); - технический шорткат:
make ddl-gpприменяет все DDL разом вручную. Команда сама не вызывается при старте контейнеров, её нужно запустить явно.
- учебный вариант: запустить DAG
- В Airflow есть коннекты:
greenplum_conn— к Greenplum;bookings_db— к сервисуbookings-db. По умолчанию они задаются через переменные окруженияAIRFLOW_CONN_...в docker-compose, поэтому могут не отображаться в UI, ноPostgresOperatorнайдёт их поconn_id. При желании их можно создать/отредактировать вручную через Admin → Connections.
3. Последовательность задач в DAG
generate_bookings_day:- PostgresOperator к
bookings-db; - выполняет SQL
/sql/src/bookings_generate_day_if_missing.sql; - скрипт смотрит на
max(book_date)вbookings.bookings:- если база пустая — берёт стартовую дату из GUC и генерирует
bookings.init_daysсуток; - если данные уже есть — добавляет один следующий учебный день после
max(book_date)и пишет NOTICE с интервалом генерации.
- если база пустая — берёт стартовую дату из GUC и генерирует
- PostgresOperator к
load_bookings_to_stg:- PostgresOperator к Greenplum;
- выполняет SQL
/sql/stg/bookings_load.sql; - считает «старый» максимум
src_created_at_ts(по предыдущим батчам) и грузит только новые строки изstg.bookings_ext, заполняяsrc_created_at_ts,load_dttm,batch_id={{ ds_nodash }}.
check_row_counts:- PostgresOperator к Greenplum;
- выполняет SQL
/sql/stg/bookings_dq.sql; - за то же окно инкремента считает количество строк в источнике и в
stg.bookings(по текущемуbatch_id); - при расхождении делает
RAISE EXCEPTIONс понятным текстом ошибки.
finish_summary:- логирует итог выполнения DAG за одно срабатывание.
4. Как этим пользоваться студенту (черновой сценарий)
- Поднять стенд и подготовить источники:
make up(Airflow инициализируется автоматически при первом старте)make bookings-initmake ddl-gp
- Открыть Airflow UI (
http://localhost:8080) и включить DAGbookings_to_gp_stage. - Вызвать
TriggerDAG (дату логического запуска можно оставить по умолчанию — она используется только как меткаbatch_id). - Посмотреть:
- в
bookings-db― что появился день с бронированиями; - в Greenplum (
make gp-psql) — данные вstg.bookings:SELECT * FROM stg.bookings LIMIT 10;SELECT src_created_at_ts, load_dttm, batch_id FROM stg.bookings ORDER BY src_created_at_ts DESC LIMIT 10;
- в
- Перезапустить DAG ещё несколько раз и увидеть, что:
- генерация в
bookings.bookingsидёт по одному дню вперёд от текущегоmax(book_date); - в
stg.bookingsпоявляются только новые записи (delta), помеченные разнымиbatch_id.
- генерация в
5. Примечания «на потом»
- Текущая документация по блоку bookings/STG разбросана:
README.md(общий обзор стенда),docs/internal/bookings_tz.md(источник bookings),docs/internal/pxf_bookings.md(PXF),docs/internal/bookings_stg_design.md(дизайн STG),- этот файл (мини‑README по DAG).
- В будущем всё это нужно будет собрать в одну понятную историю для студента:
- отдельный раздел «Учебный пример: bookings → stg → dwh»;
- скриншоты DAG, примеры запросов и типичные ошибки.