diff --git a/README.md b/README.md index a59de83..07e9313 100644 --- a/README.md +++ b/README.md @@ -80,6 +80,7 @@ make bookings-init make gp-psql -- внутри psql: SELECT COUNT(*) FROM stg.bookings; +SELECT COUNT(*) FROM stg.tickets; SELECT * FROM stg.bookings ORDER BY src_created_at_ts DESC LIMIT 10; ``` @@ -89,9 +90,9 @@ SELECT * FROM stg.bookings ORDER BY src_created_at_ts DESC LIMIT 10; Основные (для потока bookings → DWH): -- `bookings_stg_ddl` — создаёт `stg.bookings_ext` и `stg.bookings` в Greenplum; -- `bookings_to_gp_stage` — генерирует учебный день в `bookings-db` и грузит инкремент в `stg.bookings` - (через PXF), затем выполняет DQ‑проверку. +- `bookings_stg_ddl` — создаёт `stg.bookings_ext`/`stg.bookings` и `stg.tickets_ext`/`stg.tickets` в Greenplum; +- `bookings_to_gp_stage` — генерирует учебный день в `bookings-db`, грузит инкремент в `stg.bookings` и `stg.tickets` + (через PXF), затем выполняет DQ‑проверки. Вспомогательные (побочный трек с CSV): diff --git a/airflow/dags/bookings_stg_ddl.py b/airflow/dags/bookings_stg_ddl.py index 99d0715..0fc3a95 100644 --- a/airflow/dags/bookings_stg_ddl.py +++ b/airflow/dags/bookings_stg_ddl.py @@ -5,8 +5,9 @@ from __future__ import annotations Запускается вручную перед DAG загрузки bookings_to_gp_stage или после изменения DDL. """ -from datetime import datetime, timedelta +from datetime import timedelta +import pendulum from airflow.providers.postgres.operators.postgres import PostgresOperator from airflow import DAG @@ -17,7 +18,7 @@ default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(secon with DAG( dag_id="bookings_stg_ddl", - start_date=datetime(2024, 1, 1), + start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), schedule=None, catchup=False, template_searchpath="/sql", diff --git a/airflow/dags/bookings_to_gp_stage.py b/airflow/dags/bookings_to_gp_stage.py index a4982f3..8a5e21a 100644 --- a/airflow/dags/bookings_to_gp_stage.py +++ b/airflow/dags/bookings_to_gp_stage.py @@ -15,8 +15,9 @@ from __future__ import annotations - `run_id` используется как метка запуска (в `batch_id`, в логах и DQ). """ -from datetime import datetime, timedelta +from datetime import timedelta +import pendulum from airflow.operators.python import PythonOperator from airflow.providers.postgres.operators.postgres import PostgresOperator @@ -48,9 +49,10 @@ def _finish_summary() -> None: with DAG( dag_id="bookings_to_gp_stage", - start_date=datetime(2024, 1, 1), + start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), schedule=None, catchup=False, + max_active_runs=1, template_searchpath="/sql", default_args=default_args, tags=["demo", "bookings", "greenplum", "stg"], diff --git a/docs/bookings_to_gp_stage.md b/docs/bookings_to_gp_stage.md index dac1a63..e07d4e9 100644 --- a/docs/bookings_to_gp_stage.md +++ b/docs/bookings_to_gp_stage.md @@ -2,7 +2,7 @@ Этот DAG — основной учебный пример в стенде. Он показывает путь данных из источника **Postgres** (`bookings-db`, демо‑БД `demo`) в сырой слой **STG** в **Greenplum** с инкрементальной загрузкой -и простой проверкой качества данных. +и простыми проверками качества данных. ## Что делает DAG @@ -10,6 +10,8 @@ (генератор всегда “шагает” вперёд от `max(book_date)`). - В Greenplum загружает инкремент в `stg.bookings` через внешнюю таблицу `stg.bookings_ext`, используя PXF. - Сверяет количество строк между источником (за окно инкремента) и загруженным батчем. +- Загружает инкремент в `stg.tickets` через внешнюю таблицу `stg.tickets_ext`, используя PXF. +- Запускает DQ‑проверки для `stg.tickets` (количество, ссылочная целостность, обязательные поля). ## Что должно быть готово перед запуском @@ -25,9 +27,7 @@ make up make bookings-init ``` -Важно: генератор `bookings` в этом стенде поддерживается только в режиме `BOOKINGS_JOBS=1`. - -3) В Greenplum созданы `stg.bookings_ext` и `stg.bookings` (выберите один вариант): +3) В Greenplum созданы STG‑объекты `stg.bookings_ext`/`stg.bookings` и `stg.tickets_ext`/`stg.tickets` (выберите один вариант): - учебный вариант: запустить DAG `bookings_stg_ddl` в Airflow UI; - технический шорткат: `make ddl-gp`. @@ -70,7 +70,20 @@ make bookings-init вставленных в `stg.bookings` для текущего `batch_id`; - при расхождении делает `RAISE EXCEPTION` с понятным текстом. -4) `finish_summary` +4) `load_tickets_to_stg` + +- выполняет `sql/stg/tickets_load.sql` в Greenplum; +- так как в `bookings.tickets` нет явной временной колонки, окно инкремента берётся по `book_date` + из связанной внешней таблицы `stg.bookings_ext` (JOIN по `book_ref`); +- вставляет строки в `stg.tickets`, добавляя `src_created_at_ts`, `load_dttm` и `batch_id={{ run_id }}`. + +5) `check_tickets_dq` + +- выполняет `sql/stg/tickets_dq.sql` в Greenplum; +- проверяет количество строк в том же окне инкремента, а также ссылочную целостность и обязательные поля; +- при проблемах делает `RAISE EXCEPTION`, чтобы DAG падал “красным”. + +6) `finish_summary` - логирует краткую сводку в конце запуска. diff --git a/docs/chore/bookings-etl.md b/docs/chore/bookings-etl.md index 2d79659..f0bc5a7 100644 --- a/docs/chore/bookings-etl.md +++ b/docs/chore/bookings-etl.md @@ -37,27 +37,26 @@ ```sql -- sql/stg/tickets_ddl.sql -DROP TABLE IF EXISTS stg.tickets CASCADE; - -CREATE TABLE stg.tickets ( +CREATE TABLE IF NOT EXISTS stg.tickets ( -- Бизнес-атрибуты (из источника) ticket_no TEXT NOT NULL, -- номер билета book_ref TEXT NOT NULL, -- ссылка на бронирование passenger_id TEXT, -- идентификатор пассажира passenger_name TEXT, -- имя пассажира - outbound TEXT, -- флаг направления (boolean → text для Greenplum) + outbound TEXT, -- флаг направления (в источнике boolean; в STG храним как TEXT) -- Технические атрибуты src_created_at_ts TIMESTAMP, -- временная метка источника (для инкремента) - load_dttm TIMESTAMP NOT NULL, -- время загрузки в Greenplum + load_dttm TIMESTAMP NOT NULL DEFAULT now(), -- время загрузки в Greenplum batch_id TEXT -- идентификатор батча (из Airflow run_id) ) -DISTRIBUTED BY (ticket_no); -- распределение по первичному ключу +WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +DISTRIBUTED BY (book_ref); -- распределение по ключу связи с bookings ``` **Решения по проектированию:** -- **DISTRIBUTED BY (ticket_no):** равномерное распределение, так как это первичный ключ -- **boolean → TEXT:** Greenplum хорошо работает с текстовым представлением для флагов +- **DISTRIBUTED BY (book_ref):** типовой джойн `tickets → bookings` идёт по `book_ref`, так меньше motion в MPP +- **boolean → TEXT:** в сыром STG храним бизнес-колонки как TEXT (для обучения и минимизации кастов на входе) - **src_created_at_ts:** временная метка из даты связанного бронирования (см. раздел 7) ## 4. Проектирование внешней таблицы (PXF) @@ -122,7 +121,11 @@ CREATE TABLE IF NOT EXISTS stg.tickets ( batch_id TEXT ) WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) -DISTRIBUTED BY (ticket_no); +-- Распределяем по book_ref, чтобы джойны tickets → bookings по book_ref были без motion. +DISTRIBUTED BY (book_ref); + +-- На случай, если таблица уже была создана раньше с другим ключом распределения. +ALTER TABLE IF EXISTS stg.tickets SET DISTRIBUTED BY (book_ref); ``` ## 6. LOAD-скрипт (загрузка инкремента) @@ -173,11 +176,10 @@ WHERE b.book_date > COALESCE( TIMESTAMP '1900-01-01 00:00:00' ) AND NOT EXISTS ( - -- Защита от дубликатов в рамках одного батча + -- Защита от дублей: ticket_no в источнике уникален, и в stg его не дублируем. SELECT 1 FROM stg.tickets AS t - WHERE t.batch_id = '{{ run_id }}'::text - AND t.ticket_no = ext.ticket_no + WHERE t.ticket_no = ext.ticket_no ); ``` @@ -194,56 +196,83 @@ AND NOT EXISTS ( Скрипт `sql/stg/tickets_dq.sql`: ```sql --- Проверка 1: совпадение количества билетов в источнике и STG --- Используем только внешние таблицы: stg.tickets_ext + stg.bookings_ext +-- Проверки качества данных для tickets + DO $$ DECLARE - v_source_count BIGINT; - v_stg_count BIGINT; + v_batch_id TEXT := '{{ run_id }}'::text; + v_prev_ts TIMESTAMP; + v_source_count BIGINT; + v_stg_count BIGINT; + v_orphan_count BIGINT; + v_null_count BIGINT; BEGIN - -- Количество в источнике (новые билеты) - SELECT COUNT(*) INTO v_source_count + -- Опорная метка: максимум src_created_at_ts среди предыдущих батчей + SELECT max(src_created_at_ts) + INTO v_prev_ts + FROM stg.tickets + WHERE batch_id <> v_batch_id + OR batch_id IS NULL; + + -- Источник: считаем строки в том же окне инкремента, что и загрузка + SELECT COUNT(*) + INTO v_source_count FROM stg.tickets_ext AS t JOIN stg.bookings_ext AS b ON t.book_ref = b.book_ref - WHERE b.book_date > COALESCE( - (SELECT max(src_created_at_ts) FROM stg.tickets - WHERE batch_id <> '{{ run_id }}'::text OR batch_id IS NULL), - TIMESTAMP '1900-01-01 00:00:00' - ); - - -- Количество в STG (текущий батч) - SELECT COUNT(*) INTO v_stg_count + WHERE b.book_date > COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'); + + IF v_source_count = 0 THEN + RAISE EXCEPTION + 'В источнике tickets_ext нет строк для окна инкремента (book_date > %). Проверьте генерацию данных (таск generate_bookings_day).', + COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'); + END IF; + + -- STG: считаем строки текущего батча + SELECT COUNT(*) + INTO v_stg_count FROM stg.tickets - WHERE batch_id = '{{ run_id }}'::text; - - -- Проверка совпадения + WHERE batch_id = v_batch_id; + IF v_source_count <> v_stg_count THEN - RAISE EXCEPTION 'DQ FAILED: несовпадение количества билетов. Источник: %, STG: %', - v_source_count, v_stg_count; - ELSE - RAISE NOTICE 'DQ PASSED: количество билетов совпадает (%)', v_stg_count; + RAISE EXCEPTION + 'DQ FAILED: несовпадение количества билетов. Источник: %, STG (batch_id=%): %', + v_source_count, + v_batch_id, + v_stg_count; + END IF; + + -- Ссылочная целостность: tickets должны иметь соответствующие bookings в STG + SELECT COUNT(*) + INTO v_orphan_count + FROM stg.tickets AS t + WHERE t.batch_id = v_batch_id + AND NOT EXISTS ( + SELECT 1 + FROM stg.bookings AS b + WHERE b.book_ref = t.book_ref + ); + + IF v_orphan_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены tickets без соответствующих bookings (batch_id=%): %', + v_batch_id, + v_orphan_count; + END IF; + + -- Обязательные поля + SELECT COUNT(*) + INTO v_null_count + FROM stg.tickets AS t + WHERE t.batch_id = v_batch_id + AND (t.ticket_no IS NULL OR t.book_ref IS NULL); + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены tickets с NULL в обязательных полях (batch_id=%): %', + v_batch_id, + v_null_count; END IF; END $$; - --- Проверка 2: ссылочная целостность (все ticket_no должны иметь соответствующие book_ref) -SELECT - COUNT(*) AS orphan_tickets -FROM stg.tickets -WHERE batch_id = '{{ run_id }}'::text - AND NOT EXISTS ( - SELECT 1 - FROM stg.bookings - WHERE stg.bookings.book_ref = stg.tickets.book_ref - ); --- Ожидаемое значение: 0 - --- Проверка 3: отсутствие NULL в обязательных полях -SELECT - COUNT(*) AS null_tickets -FROM stg.tickets -WHERE batch_id = '{{ run_id }}'::text - AND (ticket_no IS NULL OR book_ref IS NULL); --- Ожидаемое значение: 0 ``` ## 8. Интеграция с существующими DAG @@ -396,7 +425,7 @@ WHERE b.book_ref IS NULL; 8. ✅ Проверка через Airflow UI 9. ✅ Проверка идемпотентности DDL 10. ✅ Проверка инкрементальной загрузки -11. ⏳ Обновить `README.md` (добавить tickets в список STG-таблиц) +11. ✅ Обновить `README.md` (добавить tickets в список STG-таблиц) ## 11. Сопутствующие изменения @@ -412,8 +441,8 @@ WHERE b.book_ref IS NULL; ### Необходимые изменения в документации: **README.md:** -- ⏳ Добавить `tickets` в список STG-таблиц -- ⏳ Обновить описание DAG `bookings_to_gp_stage` (упомянуть загрузку билетов) +- ✅ Добавить `tickets` в список STG-таблиц +- ✅ Обновить описание DAG `bookings_to_gp_stage` (упомянуть загрузку билетов) ## 12. Результаты тестирования @@ -431,4 +460,4 @@ WHERE b.book_ref IS NULL; ### Идемпотентность DDL: - **Первый запуск DDL:** таблицы созданы, данные не затронуты -- **Второй запуск DDL:** данные не пропали, таблицы существуют (CREATE TABLE IF NOT EXISTS) \ No newline at end of file +- **Второй запуск DDL:** данные не пропали, таблицы существуют (CREATE TABLE IF NOT EXISTS) diff --git a/sql/stg/tickets_ddl.sql b/sql/stg/tickets_ddl.sql index e43648c..553cb04 100644 --- a/sql/stg/tickets_ddl.sql +++ b/sql/stg/tickets_ddl.sql @@ -29,4 +29,6 @@ CREATE TABLE IF NOT EXISTS stg.tickets ( batch_id TEXT ) WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) -DISTRIBUTED BY (ticket_no); +-- Распределяем по book_ref, чтобы джойны tickets → bookings по book_ref были без motion. +DISTRIBUTED BY (book_ref); + diff --git a/sql/stg/tickets_dq.sql b/sql/stg/tickets_dq.sql index b5f5feb..64f63e9 100644 --- a/sql/stg/tickets_dq.sql +++ b/sql/stg/tickets_dq.sql @@ -1,52 +1,83 @@ -- Проверки качества данных для tickets --- Проверка 1: совпадение количества билетов в источнике и STG DO $$ DECLARE - v_source_count BIGINT; - v_stg_count BIGINT; + v_batch_id TEXT := '{{ run_id }}'::text; + v_prev_ts TIMESTAMP; + v_source_count BIGINT; + v_stg_count BIGINT; + v_orphan_count BIGINT; + v_null_count BIGINT; BEGIN - -- Количество в источнике (новые билеты) - -- Используем только внешние таблицы: stg.tickets_ext + stg.bookings_ext - SELECT COUNT(*) INTO v_source_count + -- Опорная метка: максимум src_created_at_ts среди предыдущих батчей + SELECT max(src_created_at_ts) + INTO v_prev_ts + FROM stg.tickets + WHERE batch_id <> v_batch_id + OR batch_id IS NULL; + + -- Количество в источнике (новые билеты в том же окне инкремента, что и загрузка) + SELECT COUNT(*) + INTO v_source_count FROM stg.tickets_ext AS t JOIN stg.bookings_ext AS b ON t.book_ref = b.book_ref - WHERE b.book_date > COALESCE( - (SELECT max(src_created_at_ts) FROM stg.tickets - WHERE batch_id <> '{{ run_id }}'::text OR batch_id IS NULL), - TIMESTAMP '1900-01-01 00:00:00' - ); + WHERE b.book_date > COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'); + + IF v_source_count = 0 THEN + RAISE EXCEPTION + 'В источнике tickets_ext нет строк для окна инкремента (book_date > %). Проверьте генерацию данных (таск generate_bookings_day).', + COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'); + END IF; -- Количество в STG (текущий батч) - SELECT COUNT(*) INTO v_stg_count + SELECT COUNT(*) + INTO v_stg_count FROM stg.tickets - WHERE batch_id = '{{ run_id }}'::text; + WHERE batch_id = v_batch_id; - -- Проверка совпадения IF v_source_count <> v_stg_count THEN - RAISE EXCEPTION 'DQ FAILED: несовпадение количества билетов. Источник: %, STG: %', - v_source_count, v_stg_count; - ELSE - RAISE NOTICE 'DQ PASSED: количество билетов совпадает (%)', v_stg_count; + RAISE EXCEPTION + 'DQ FAILED: несовпадение количества билетов. Источник: %, STG (batch_id=%): %', + v_source_count, + v_batch_id, + v_stg_count; END IF; + + -- Проверка ссылочной целостности: все tickets должны иметь соответствующие bookings в этом же STG + SELECT COUNT(*) + INTO v_orphan_count + FROM stg.tickets AS t + WHERE t.batch_id = v_batch_id + AND NOT EXISTS ( + SELECT 1 + FROM stg.bookings AS b + WHERE b.book_ref = t.book_ref + ); + + IF v_orphan_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены tickets без соответствующих bookings (batch_id=%): %', + v_batch_id, + v_orphan_count; + END IF; + + -- Проверка обязательных полей + SELECT COUNT(*) + INTO v_null_count + FROM stg.tickets AS t + WHERE t.batch_id = v_batch_id + AND (t.ticket_no IS NULL OR t.book_ref IS NULL); + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены tickets с NULL в обязательных полях (batch_id=%): %', + v_batch_id, + v_null_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: tickets ок (batch_id=%): source=% stg=%', + v_batch_id, + v_source_count, + v_stg_count; END $$; - --- Проверка 2: ссылочная целостность (все ticket_no должны иметь соответствующие book_ref) -SELECT - COUNT(*) AS orphan_tickets -FROM stg.tickets -WHERE batch_id = '{{ run_id }}'::text - AND NOT EXISTS ( - SELECT 1 - FROM stg.bookings - WHERE stg.bookings.book_ref = stg.tickets.book_ref - ); --- Ожидаемое значение: 0 - --- Проверка 3: отсутствие NULL в обязательных полях -SELECT - COUNT(*) AS null_tickets -FROM stg.tickets -WHERE batch_id = '{{ run_id }}'::text - AND (ticket_no IS NULL OR book_ref IS NULL); --- Ожидаемое значение: 0 diff --git a/sql/stg/tickets_load.sql b/sql/stg/tickets_load.sql index 73a3ad5..3a32059 100644 --- a/sql/stg/tickets_load.sql +++ b/sql/stg/tickets_load.sql @@ -32,9 +32,8 @@ WHERE b.book_date > COALESCE( TIMESTAMP '1900-01-01 00:00:00' ) AND NOT EXISTS ( - -- Защита от дубликатов в рамках одного батча + -- Защита от дублей: ticket_no в источнике уникален, и в stg его не дублируем. SELECT 1 FROM stg.tickets AS t - WHERE t.batch_id = '{{ run_id }}'::text - AND t.ticket_no = ext.ticket_no + WHERE t.ticket_no = ext.ticket_no );