From 29490828f45389a7bff5b178e5c8a2e2d96773cd Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Wed, 10 Dec 2025 21:53:17 +0300 Subject: [PATCH] =?UTF-8?q?=D0=A3=D1=85=D0=BE=D0=B4=20=D0=BE=D1=82=20?= =?UTF-8?q?=D0=B7=D0=B0=D0=BF=D1=83=D1=81=D0=BA=D0=B0=20DAG=20=D0=B7=D0=B0?= =?UTF-8?q?=20=D0=BA=D0=BE=D0=BD=D0=BA=D1=80=D0=B5=D1=82=D0=BD=D1=83=D1=8E?= =?UTF-8?q?=20=D0=B4=D0=B0=D1=82=D1=83.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Теперь новый запуск генерирует и переливает новый, следующий день --- README.md | 17 +++++++++---- TESTING.md | 7 +++--- airflow/dags/bookings_to_gp_stage.py | 14 +++++++---- docs/internal/bookings_stg_design.md | 16 ++++++------- docs/internal/bookings_stg_readme.md | 13 +++++----- educational-tasks.md | 19 +++++++++++++++ sql/src/bookings_generate_day_if_missing.sql | 25 +++++++------------- sql/stg/bookings_dq.sql | 23 +++++++----------- sql/stg/bookings_load.sql | 8 +++---- 9 files changed, 80 insertions(+), 62 deletions(-) diff --git a/README.md b/README.md index ad30171..7da2a6f 100644 --- a/README.md +++ b/README.md @@ -199,6 +199,11 @@ docker compose -f docker-compose.yml exec bookings-db bash -lc 'PGPASSWORD="$POS ``` ### Генерация следующего дня в bookings + +- При `make bookings-init` автоматически генерируется `BOOKINGS_INIT_DAYS` суток, начиная с даты `BOOKINGS_START_DATE` (по умолчанию один день с `2017-01-01`). +- Дальше каждый вызов `make bookings-generate-day` или генерации через DAG добавляет ровно **один** следующий день после `max(book_date)` в `bookings.bookings` — генератор сам смотрит последнюю дату. +- Рекомендуемый учебный сценарий для DAG `bookings_to_gp_stage`: запускать DAG по одному дню вперёд, выбирая в форме Trigger логическую дату `Execution Date (ds)`, совпадающую с тем днём, который вы хотите загрузить (например, `2017-01-01`, затем `2017-01-02` и т.д.). + - Быстрее всего: `make bookings-generate-day` — читает GUC и сам вызывает `continue`. - Вручную из psql/DBeaver: ```sql @@ -269,10 +274,10 @@ docker compose -f docker-compose.yml exec bookings-db bash -lc 'PGPASSWORD="$POS - подключение к БД через Airflow Connections (`bookings_db`, `greenplum_conn`); - бизнес-логика инкрементальной загрузки и DQ вынесена в SQL-файлы в каталоге `sql/`: - - `sql/src/bookings_generate_day_if_missing.sql` — генерация учебного дня в демо-БД bookings (идемпотентно); + - `sql/src/bookings_generate_day_if_missing.sql` — генерация следующего учебного дня в демо-БД bookings (или нескольких стартовых дней, если база пуста); - `sql/stg/bookings_ddl.sql` — DDL для схемы `stg` и таблиц `stg.bookings_ext` / `stg.bookings`; - - `sql/stg/bookings_load.sql` — загрузка инкремента из `stg.bookings_ext` в `stg.bookings`; - - `sql/stg/bookings_dq.sql` — проверка количества строк между источником и stg. + - `sql/stg/bookings_load.sql` — загрузка инкремента из `stg.bookings_ext` в `stg.bookings` на основе «хвоста» после предыдущих батчей; + - `sql/stg/bookings_dq.sql` — проверка количества строк между источником и stg за то же окно. Фрагмент DAG: @@ -281,7 +286,7 @@ load_bookings_to_stg = PostgresOperator( task_id="load_bookings_to_stg", postgres_conn_id="greenplum_conn", sql="/sql/stg/bookings_load.sql", - params={"load_date": "{{ ds }}", "batch_id": "{{ ds_nodash }}"}, + params={"batch_id": "{{ ds_nodash }}"}, ) ``` @@ -353,9 +358,13 @@ load_bookings_to_stg = PostgresOperator( |----------|---------| | Airflow UI не открывается | Дождитесь сообщения `Listening at: http://0.0.0.0:8080` в логах (`make logs`) | | Ошибка подключения к Greenplum | Убедитесь, что контейнер `greenplum` стал статусом `healthy` (проверьте `docker compose ps`) | +| Не открывается порт 8080/5433/5434/5435 | Проверьте, что эти порты не заняты локальными сервисами; при необходимости остановите их или измените порты в `.env`/`docker-compose.yml` | | Нет файла в `./data` после запуска DAG | Проверьте логи задачи `generate_csv`, убедитесь, что `CSV_DIR` смонтирован в docker-compose | | Команда `make` не найдена | Используйте полные команды `docker compose` или установите make | | Greenplum не стартует/падает при старте | Выполните `make down`, затем `make up && make airflow-init` (очищает тома и поднимает заново) | +| DAG `bookings_to_gp_stage` падает на внешней таблице/подключении к bookings | Убедитесь, что запущен контейнер `bookings-db` (`docker compose ps`, при необходимости `docker compose start bookings-db`), и выполнены `make bookings-init` и `make ddl-gp` или DAG `bookings_stg_ddl` | +| DAG `bookings_to_gp_stage` ругается на отсутствующие таблицы stg | Запустите DAG `bookings_stg_ddl` (или выполните `make ddl-gp`), затем повторите запуск | +| DAG не видит Greenplum/DEMObase по Airflow Connections | Проверьте, что в Airflow созданы подключения `greenplum_conn` и `bookings_db` с параметрами из раздела «Настройка подключения к Greenplum в Airflow» и блока про bookings | --- diff --git a/TESTING.md b/TESTING.md index 56986e3..a8ec617 100644 --- a/TESTING.md +++ b/TESTING.md @@ -12,8 +12,8 @@ ## 2. Локальные автоматические проверки (без Docker) - `make test` — короткие unit-тесты (`tests/test_greenplum_helpers.py`, `tests/test_dags_smoke.py`). - Smoke-тесты DAG автоматически `skip`, если Airflow не установлен в venv, поэтому прогонится за миллисекунды. -- `make lint` — black/isort в режиме проверки. Сейчас упадёт из‑за форматирования DAG-файлов. -- `make fmt` — автоисправление форматирования; после этого `make lint` должен пройти. +- `make lint` — black/isort в режиме проверки (после `make fmt` должен проходить без ошибок). +- `make fmt` — автоисправление форматирования; полезно запускать перед пушем. - (опционально) `uv run pytest -q -k dags_smoke` — только DAG smoke. ## 3. Подготовка Docker-стенда @@ -51,6 +51,7 @@ - **Проблемы с подключением**: временно изменить `GP_HOST` или `GP_PORT` на несуществующий, перезапустить `make up`, убедиться, что DAG падает с понятной ошибкой (`psycopg2.OperationalError`). - **Fallback без Airflow Connection**: установить `GP_USE_AIRFLOW_CONN=false`, перезапустить стек (`make down && make up && make airflow-init`), удостовериться, что загрузка и DQ работают через ENV. - **Дубликаты**: дважды вызвать `csv_to_greenplum` — ожидаем, что количество строк в `public.orders` не увеличится на размер CSV, а DAG `csv_to_greenplum_dq` не найдёт дублей. +- **PXF и демобаза bookings** (после настройки PXF и выполнения `make ddl-gp`): временно остановить `bookings-db` (`docker compose stop bookings-db`) и попробовать выполнить `SELECT COUNT(*) FROM public.ext_bookings_bookings;` в `make gp-psql` — ожидается ошибка подключения. Затем запустить `bookings-db` (`docker compose start bookings-db`) и убедиться, что запрос снова работает. ## 7. Быстрый reset (если «что-то сломалось») - Перезапустить стенд с очисткой данных: @@ -69,5 +70,5 @@ ## Текущий статус (пример успешного прогона) - `uv run pytest -q` — 11 passed, 2 smoke-теста DAG пропущены (Airflow не установлен в venv). -- `make lint` — падает, потому что `airflow/dags/*.py` не отформатированы black/isort. После `make fmt` проблема уйдёт. +- `make lint` — проходит (DAG‑файлы отформатированы black/isort). - Docker-стенд не запускался в рамках этой сессии; ожидается, что инструкции выше обеспечат полноценную проверку. diff --git a/airflow/dags/bookings_to_gp_stage.py b/airflow/dags/bookings_to_gp_stage.py index de4df21..7d6a1b2 100644 --- a/airflow/dags/bookings_to_gp_stage.py +++ b/airflow/dags/bookings_to_gp_stage.py @@ -8,6 +8,15 @@ from __future__ import annotations - Airflow Connections для подключения к БД; - PostgresOperator, который берёт SQL-скрипты с диска; чтобы студент увидел «канонический» способ работы с SQL в DAG. + +Каждый запуск DAG работает как «шаг по времени вперёд»: +- генератор в демо-БД bookings добавляет следующий учебный день после max(book_date); +- загрузка в Greenplum берёт все строки, появившиеся после предыдущих батчей; +- `batch_id` (через `ds_nodash`) помечает строки конкретного запуска. + +Логическая дата запуска DAG (`ds`) используется только как удобная метка +запуска (для `batch_id`, логов и DQ), но не управляет тем, за какие дни +генерируются и переливаются данные. """ from datetime import datetime, timedelta @@ -50,12 +59,11 @@ with DAG( tags=["demo", "bookings", "greenplum", "stg"], description="Учебный DAG: загрузка из bookings-db в stg.bookings (Greenplum)", ) as dag: - # 1. Генерируем один учебный день в демо-БД bookings (идемпотентно) + # 1. Генерируем один (или несколько стартовых) учебный день в демо-БД bookings generate_bookings_day = PostgresOperator( task_id="generate_bookings_day", postgres_conn_id=BOOKINGS_CONN_ID, sql="/sql/src/bookings_generate_day_if_missing.sql", - params={"load_date": "{{ ds }}"}, ) # 2. Загружаем инкремент из stg.bookings_ext в stg.bookings @@ -64,7 +72,6 @@ with DAG( postgres_conn_id=GREENPLUM_CONN_ID, sql="/sql/stg/bookings_load.sql", params={ - "load_date": "{{ ds }}", "batch_id": "{{ ds_nodash }}", }, ) @@ -75,7 +82,6 @@ with DAG( postgres_conn_id=GREENPLUM_CONN_ID, sql="/sql/stg/bookings_dq.sql", params={ - "load_date": "{{ ds }}", "batch_id": "{{ ds_nodash }}", }, ) diff --git a/docs/internal/bookings_stg_design.md b/docs/internal/bookings_stg_design.md index 3dee531..cc928e4 100644 --- a/docs/internal/bookings_stg_design.md +++ b/docs/internal/bookings_stg_design.md @@ -84,11 +84,11 @@ DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Gre - Этот DAG не загружает данные, только подготавливает структуру. - Вся DDL‑логика (CREATE/ALTER/DROP) сосредоточена здесь; рабочие DAG’и занимаются только DML (INSERT/SELECT). -### 4.2. DAG для ежедневной загрузки +### 4.2. DAG для пошаговой загрузки - `dag_id`: `bookings_to_gp_stage`. - Основные параметры: - - `load_date` (по умолчанию `{{ ds }}`) — учебный день, за который генерим и грузим данные; + - `batch_id` (по умолчанию `{{ ds_nodash }}`) — метка батча, которая попадает в `stg.bookings.batch_id`; - подключения: - `bookings_db_conn_id` — Airflow connection к `bookings-db` (в коде DAG — `BOOKINGS_CONN_ID = "bookings_db"`); - `greenplum_conn_id` — Airflow connection к Greenplum (`GREENPLUM_CONN_ID = "greenplum_conn"`). @@ -97,22 +97,22 @@ DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Gre 1. `generate_bookings_day` - PostgresOperator к `bookings-db`; - - выполняет скрипт `/sql/bookings/generate_day_if_missing.sql`; - - скрипт сам идемпотентно проверяет, есть ли данные за `load_date` в `bookings.bookings`: - - если день уже есть — выдаёт NOTICE и ничего не делает; - - если нет — вызывает генератор демобазы (логика как в `bookings/generate_next_day.sql`). + - выполняет скрипт `/sql/src/bookings_generate_day_if_missing.sql`; + - скрипт смотрит на `max(book_date)` и: + - если база пуста — берёт стартовую дату из конфигурации (`bookings.start_date`) и генерирует `bookings.init_days` суток; + - если данные уже есть — добавляет один следующий учебный день после `max(book_date)` (логика как в `bookings/generate_next_day.sql`). 2. `load_bookings_to_stg` - PostgresOperator к Greenplum; - выполняет скрипт `/sql/stg/bookings_load.sql`; - внутри SQL считается `max(src_created_at_ts)` по «старым» батчам и по нему строится окно инкремента: - первая загрузка (full) — берём все строки из `stg.bookings_ext`; - - последующие загрузки — берём только записи, где `book_date` больше предыдущего максимума и не позже конца учебного дня; + - последующие загрузки — берём только записи, где `book_date` больше предыдущего максимума (верхняя граница по дате не задаётся явно); - при вставке заполняются тех.колонки `src_created_at_ts`, `load_dttm`, `batch_id`. 3. `check_row_counts` - PostgresOperator к Greenplum; - выполняет скрипт `/sql/stg/bookings_dq.sql`; - скрипт заново считает окно инкремента по тем же правилам, что и загрузка, и сравнивает: - - количество строк в `stg.bookings_ext` за окно, + - количество строк в `stg.bookings_ext` с `book_date` позже «старого» максимума, - количество строк в `stg.bookings` для текущего `batch_id`; - при расхождении выполняет `RAISE EXCEPTION` с понятным текстом ошибки. 4. `finish_summary` diff --git a/docs/internal/bookings_stg_readme.md b/docs/internal/bookings_stg_readme.md index 8ab8dc4..76e9318 100644 --- a/docs/internal/bookings_stg_readme.md +++ b/docs/internal/bookings_stg_readme.md @@ -32,8 +32,9 @@ _Внутренний файл, чтобы не забыть договорён - `generate_bookings_day`: - PostgresOperator к `bookings-db`; - выполняет SQL `/sql/src/bookings_generate_day_if_missing.sql`; - - скрипт сам проверяет, есть ли строки за `load_date` (`book_date::date = load_date`); - - если день уже сгенерирован — ничего не делает (идемпотентность), только пишет NOTICE. + - скрипт смотрит на `max(book_date)` в `bookings.bookings`: + - если база пустая — берёт стартовую дату из GUC и генерирует `bookings.init_days` суток; + - если данные уже есть — добавляет один следующий учебный день после `max(book_date)` и пишет NOTICE с интервалом генерации. - `load_bookings_to_stg`: - PostgresOperator к Greenplum; - выполняет SQL `/sql/stg/bookings_load.sql`; @@ -54,15 +55,15 @@ _Внутренний файл, чтобы не забыть договорён - `make bookings-init` - `make ddl-gp` 2. Открыть Airflow UI (`http://localhost:8080`) и включить DAG `bookings_to_gp_stage`. -3. Вызвать `Trigger` DAG с конкретной датой (например, `2017-01-01` → зависит от стартовой конфигурации демобазы). +3. Вызвать `Trigger` DAG (дату логического запуска можно оставить по умолчанию — она используется только как метка `batch_id`). 4. Посмотреть: - в `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;` -5. Перезапустить DAG для следующей даты и увидеть, что: - - генерация в `bookings.bookings` идёт по одному дню вперёд; - - в `stg.bookings` появляются только новые записи (delta). +5. Перезапустить DAG ещё несколько раз и увидеть, что: + - генерация в `bookings.bookings` идёт по одному дню вперёд от текущего `max(book_date)`; + - в `stg.bookings` появляются только новые записи (delta), помеченные разными `batch_id`. ## 5. Примечания «на потом» diff --git a/educational-tasks.md b/educational-tasks.md index 4e9c28b..be2792c 100644 --- a/educational-tasks.md +++ b/educational-tasks.md @@ -81,6 +81,22 @@ - `\dn` и `\dt stg.*` - `SELECT * FROM stg.bookings LIMIT 5;` (после запуска соответствующего DAG). +### 2.3. Как генерируются учебные данные bookings + +1. Откройте файл `bookings/generate_next_day.sql` и ответьте себе на вопросы: + - с какой даты начинается генерация данных (посмотрите на GUC `bookings.start_date` и переменную `v_start_cfg`); + - сколько дней генерируется при первой установке (переменная `bookings.init_days`); + - что происходит, если таблица `bookings.bookings` уже не пустая. +2. В демобазе (`make bookings-psql`) выполните: + - `SELECT min(book_date), max(book_date) FROM bookings.bookings;` + - затем запустите `make bookings-generate-day` и повторите запрос — как изменился максимальный день? +3. Откройте `sql/src/bookings_generate_day_if_missing.sql` и обратите внимание, что: + - логическая дата запуска DAG (`{{ ds }}`) не влияет на выбор дня генерации; + - скрипт всегда смотрит на `max(book_date)` и добавляет **следующий** день (или несколько стартовых дней, если база пуста). +4. Сделайте вывод: генератор всегда «шагает» по датам вперёд от максимальной даты, поэтому: + - при `make bookings-init` вы получаете `BOOKINGS_INIT_DAYS` дней начиная с `BOOKINGS_START_DATE`; + - при последующих вызовах (`make bookings-generate-day` или DAG) добавляется ровно один новый день. + --- ## 3. DAG bookings_to_gp_stage (заготовка заданий) @@ -99,6 +115,9 @@ - генерация учебного дня в `bookings.bookings`; - загрузка инкремента в `stg.bookings`; - проверка количества строк между источником и STG. +4. Обратите внимание, как в DAG используется логическая дата запуска: + - `{{ ds_nodash }}` используется как `batch_id` — метка загрузки в таблице `stg.bookings` для конкретного запуска; + - сами даты данных (какие дни есть в `bookings.bookings`) определяются генератором по `max(book_date)`, а не по `ds`. На этом этапе достаточно понять общую цепочку. Детальные задания по переработке модели данных и построению ODS/DDS/DM слоёв будут добавлены позже. diff --git a/sql/src/bookings_generate_day_if_missing.sql b/sql/src/bookings_generate_day_if_missing.sql index d4b7692..fed1b36 100644 --- a/sql/src/bookings_generate_day_if_missing.sql +++ b/sql/src/bookings_generate_day_if_missing.sql @@ -1,10 +1,11 @@ -- Генерация учебного дня в демо-БД bookings. --- Если данные за указанный день уже существуют, генерация не выполняется. +-- Этот скрипт всегда добавляет следующий учебный день после максимальной даты +-- в таблице bookings.bookings (или несколько дней от стартовой даты, если база пуста). +-- Логическая дата запуска DAG ({{ ds }}) здесь не используется для выбора дня — +-- она служит только меткой батча и попадает в логи/стейдж. DO $$ DECLARE - v_has_day boolean; - v_load_date date := {{ params.load_date }}::date; v_max_book_date timestamptz; v_start_date timestamptz; v_end_date timestamptz; @@ -17,19 +18,6 @@ BEGIN RAISE EXCEPTION 'Таблица bookings.bookings не найдена. Сначала выполните make bookings-init.'; END IF; - -- Проверяем, есть ли уже данные за нужный день - SELECT EXISTS ( - SELECT 1 - FROM bookings.bookings - WHERE book_date::date = v_load_date - ) - INTO v_has_day; - - IF v_has_day THEN - RAISE NOTICE 'Данные за % уже есть в bookings.bookings — генерация не требуется.', v_load_date; - RETURN; - END IF; - -- Ищем последнюю сгенерированную дату SELECT max(book_date) INTO v_max_book_date FROM bookings.bookings; @@ -50,10 +38,13 @@ BEGIN CALL continue(v_end_date, v_jobs); END IF; + RAISE NOTICE 'Сгенерированы данные в bookings.bookings за интервал [% - %).', + date_trunc('day', v_start_date), + date_trunc('day', v_end_date); + -- Ждём завершения фоновых джобов генератора, чтобы данные успели записаться WHILE busy() LOOP PERFORM pg_sleep(1); END LOOP; PERFORM dblink_disconnect(unnest(dblink_get_connections())); END $$; - diff --git a/sql/stg/bookings_dq.sql b/sql/stg/bookings_dq.sql index 879faf3..faa54ca 100644 --- a/sql/stg/bookings_dq.sql +++ b/sql/stg/bookings_dq.sql @@ -1,10 +1,10 @@ -- Проверка количества строк между источником stg.bookings_ext и стейджем stg.bookings. --- Считаем строки за то же окно инкремента, что и при загрузке, --- используя batch_id и "старый" максимум src_created_at_ts. +-- Считаем строки за то же окно инкремента, что и при загрузке: +-- все записи во внешней таблице с book_date больше максимального src_created_at_ts +-- из предыдущих батчей должны совпасть по количеству со строками текущего batch_id. DO $$ DECLARE - v_load_date date := {{ params.load_date }}::date; v_batch_id text := {{ params.batch_id | tojson }}::text; v_prev_ts timestamp; v_src_count bigint; @@ -17,17 +17,11 @@ BEGIN WHERE batch_id <> v_batch_id OR batch_id IS NULL; - IF v_prev_ts IS NULL THEN - -- Полная загрузка: считаем все строки во внешней таблице - SELECT COUNT(*) INTO v_src_count FROM stg.bookings_ext; - ELSE - -- Инкремент: считаем только строки за текущее окно - SELECT COUNT(*) - INTO v_src_count - FROM stg.bookings_ext - WHERE book_date > v_prev_ts - AND book_date <= (v_load_date + INTERVAL '1 day'); - END IF; + -- Источник: считаем строки во внешней таблице, которые вошли в новое окно + SELECT COUNT(*) + INTO v_src_count + FROM stg.bookings_ext + WHERE book_date > COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'); -- Считаем строки, реально вставленные в stg.bookings в этом батче SELECT COUNT(*) @@ -47,4 +41,3 @@ BEGIN v_src_count, v_stg_count; END $$; - diff --git a/sql/stg/bookings_load.sql b/sql/stg/bookings_load.sql index ea2908f..2dbef44 100644 --- a/sql/stg/bookings_load.sql +++ b/sql/stg/bookings_load.sql @@ -1,7 +1,7 @@ -- Загрузка инкремента из stg.bookings_ext в stg.bookings. -- Окно инкремента определяется по src_created_at_ts: --- берем строки, где book_date больше максимального src_created_at_ts --- среди "старых" батчей и не позже конца учебного дня. +-- берём строки, где book_date больше максимального src_created_at_ts +-- среди "старых" батчей; верхняя граница по дате не используется. INSERT INTO stg.bookings ( book_ref, @@ -27,6 +27,4 @@ WHERE book_date > COALESCE( OR batch_id IS NULL ), TIMESTAMP '1900-01-01 00:00:00' -) - AND book_date <= ({{ params.load_date }}::date + INTERVAL '1 day'); - +);