From 897a76588d7e6fb3c0c5360848116ceb4498a440 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 17 Jan 2026 15:09:08 +0300 Subject: [PATCH] =?UTF-8?q?=D0=9F=D0=B5=D1=80=D0=B2=D0=B0=D1=8F=20=D0=B2?= =?UTF-8?q?=D0=B5=D1=80=D1=81=D0=B8=D1=8F=20dag?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- airflow/dags/bookings_to_gp_stage.py | 21 +++++++++-- sql/stg/tickets_ddl.sql | 31 ++++++++++++++++ sql/stg/tickets_dq.sql | 54 ++++++++++++++++++++++++++++ sql/stg/tickets_load.sql | 49 +++++++++++++++++++++++++ 4 files changed, 152 insertions(+), 3 deletions(-) create mode 100644 sql/stg/tickets_ddl.sql create mode 100644 sql/stg/tickets_dq.sql create mode 100644 sql/stg/tickets_load.sql diff --git a/airflow/dags/bookings_to_gp_stage.py b/airflow/dags/bookings_to_gp_stage.py index 1112106..a4982f3 100644 --- a/airflow/dags/bookings_to_gp_stage.py +++ b/airflow/dags/bookings_to_gp_stage.py @@ -54,7 +54,7 @@ with DAG( template_searchpath="/sql", default_args=default_args, tags=["demo", "bookings", "greenplum", "stg"], - description="Учебный DAG: загрузка из bookings-db в stg.bookings (Greenplum)", + description="Учебный DAG: загрузка из bookings-db в stg.bookings и stg.tickets (Greenplum)", ) as dag: # 1. Генерируем один (или несколько стартовых) учебный день в демо-БД bookings generate_bookings_day = PostgresOperator( @@ -78,10 +78,25 @@ with DAG( sql="stg/bookings_dq.sql", ) - # 4. Финальный лог/сводка + # 4. Загружаем инкремент билетов из stg.tickets_ext в stg.tickets + load_tickets_to_stg = PostgresOperator( + task_id="load_tickets_to_stg", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/tickets_load.sql", + ) + + # 5. Проверяем качество данных для tickets + check_tickets_dq = PostgresOperator( + task_id="check_tickets_dq", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/tickets_dq.sql", + ) + + # 6. Финальный лог/сводка finish_summary = PythonOperator( task_id="finish_summary", python_callable=_finish_summary, ) - generate_bookings_day >> load_bookings_to_stg >> check_row_counts >> finish_summary + generate_bookings_day >> load_bookings_to_stg >> check_row_counts + check_row_counts >> load_tickets_to_stg >> check_tickets_dq >> finish_summary diff --git a/sql/stg/tickets_ddl.sql b/sql/stg/tickets_ddl.sql new file mode 100644 index 0000000..bb101cb --- /dev/null +++ b/sql/stg/tickets_ddl.sql @@ -0,0 +1,31 @@ +-- Создание внешней таблицы для доступа к bookings.tickets через PXF + +DROP EXTERNAL TABLE IF EXISTS stg.tickets_ext CASCADE; + +CREATE EXTERNAL TABLE stg.tickets_ext ( + ticket_no TEXT, + book_ref TEXT, + passenger_id TEXT, + passenger_name TEXT, + outbound TEXT +) +LOCATION ('pxf://bookings-db:5432/demo?PROFILE=postgres&SERVER=bookings_db') +FORMAT 'CUSTOM' (FORMATTER='pxfwritable_import') +ENCODING 'UTF8'; + +-- Создание внутренней таблицы для хранения данных в Greenplum + +DROP TABLE IF EXISTS stg.tickets CASCADE; + +CREATE TABLE stg.tickets ( + ticket_no TEXT NOT NULL, + book_ref TEXT NOT NULL, + passenger_id TEXT, + passenger_name TEXT, + outbound TEXT, + + src_created_at_ts TIMESTAMP, + load_dttm TIMESTAMP NOT NULL, + batch_id TEXT +) +DISTRIBUTED BY (ticket_no); diff --git a/sql/stg/tickets_dq.sql b/sql/stg/tickets_dq.sql new file mode 100644 index 0000000..d0db1af --- /dev/null +++ b/sql/stg/tickets_dq.sql @@ -0,0 +1,54 @@ +-- Проверки качества данных для tickets + +-- Проверка 1: совпадение количества билетов в источнике и STG +DO $$ +DECLARE + v_source_count BIGINT; + v_stg_count BIGINT; +BEGIN + -- Количество в источнике (новые билеты) + SELECT COUNT(*) INTO v_source_count + FROM ( + SELECT t.ticket_no + 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' + ) + ) AS source; + + -- Количество в STG (текущий батч) + SELECT COUNT(*) INTO v_stg_count + FROM stg.tickets + WHERE batch_id = '{{ run_id }}'::text; + + -- Проверка совпадения + 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; + 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 diff --git a/sql/stg/tickets_load.sql b/sql/stg/tickets_load.sql new file mode 100644 index 0000000..44b9627 --- /dev/null +++ b/sql/stg/tickets_load.sql @@ -0,0 +1,49 @@ +-- Загрузка инкремента из stg.tickets_ext в stg.tickets +-- Инкремент определяется по дате бронирования (book_date из bookings.bookings) + +INSERT INTO stg.tickets ( + ticket_no, + book_ref, + passenger_id, + passenger_name, + outbound, + src_created_at_ts, + load_dttm, + batch_id +) +SELECT + ext.ticket_no, + ext.book_ref, + ext.passenger_id, + ext.passenger_name, + ext.outbound, + b.book_date::timestamp, -- временная метка из бронирования + now(), + '{{ run_id }}'::text +FROM stg.tickets_ext AS ext +JOIN ( + -- Подзапрос: получаем book_date для инкремента + -- Связываем tickets с bookings через внешнюю таблицу stg.bookings_ext + SELECT + b.book_ref, + b.book_date + FROM stg.bookings_ext AS b_ext + JOIN bookings.bookings AS b ON b_ext.book_ref = b.book_ref + -- Берём только новые бронирования (по дате) + WHERE b_ext.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' + ) +) AS b ON ext.book_ref = b.book_ref +AND NOT EXISTS ( + -- Защита от дубликатов в рамках одного батча + SELECT 1 + FROM stg.tickets AS t + WHERE t.batch_id = '{{ run_id }}'::text + AND t.ticket_no = ext.ticket_no +);