diff --git a/airflow/dags/bookings_stg_ddl.py b/airflow/dags/bookings_stg_ddl.py index 7c17cdd..99d0715 100644 --- a/airflow/dags/bookings_stg_ddl.py +++ b/airflow/dags/bookings_stg_ddl.py @@ -1,7 +1,7 @@ from __future__ import annotations """ -Учебный DAG: создаёт схему stg и таблицы bookings_ext/bookings в Greenplum. +Учебный DAG: создаёт схему stg и таблицы bookings_ext/bookings/tickets в Greenplum. Запускается вручную перед DAG загрузки bookings_to_gp_stage или после изменения DDL. """ @@ -22,11 +22,19 @@ with DAG( catchup=False, template_searchpath="/sql", default_args=default_args, - tags=["demo", "greenplum", "ddl", "bookings", "stg"], - description="Создаёт/обновляет stg.bookings_ext и stg.bookings для учебного DAG", + tags=["demo", "greenplum", "ddl", "bookings", "tickets", "stg"], + description="Создаёт/обновляет stg.bookings_ext/bookings/tickets для учебного DAG", ) as dag: apply_stg_bookings_ddl = PostgresOperator( task_id="apply_stg_bookings_ddl", postgres_conn_id=GREENPLUM_CONN_ID, sql="stg/bookings_ddl.sql", ) + + apply_stg_tickets_ddl = PostgresOperator( + task_id="apply_stg_tickets_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/tickets_ddl.sql", + ) + + apply_stg_bookings_ddl >> apply_stg_tickets_ddl diff --git a/sql/ddl_gp.sql b/sql/ddl_gp.sql index 0e81d60..e95f0e7 100644 --- a/sql/ddl_gp.sql +++ b/sql/ddl_gp.sql @@ -22,6 +22,7 @@ CREATE EXTERNAL TABLE public.ext_bookings_bookings ( LOCATION ('pxf://bookings.bookings?PROFILE=JDBC&SERVER=bookings-db') FORMAT 'CUSTOM' (formatter='pxfwritable_import'); --- DDL для слоя stg по таблице bookings вынесен в отдельный файл. --- Здесь подключаем его через psql \i, чтобы сохранить единый входной скрипт. +-- DDL для слоя stg по таблицам bookings и tickets вынесены в отдельные файлы. +-- Здесь подключаем их через psql \i, чтобы сохранить единый входной скрипт. \i stg/bookings_ddl.sql +\i stg/tickets_ddl.sql diff --git a/sql/stg/tickets_ddl.sql b/sql/stg/tickets_ddl.sql index bb101cb..e43648c 100644 --- a/sql/stg/tickets_ddl.sql +++ b/sql/stg/tickets_ddl.sql @@ -1,7 +1,12 @@ --- Создание внешней таблицы для доступа к bookings.tickets через PXF +-- DDL для слоя STG по таблице tickets. +-- Используется как из общего скрипта ddl_gp.sql (через \i), +-- так и может выполняться отдельно при изменении схемы. -DROP EXTERNAL TABLE IF EXISTS stg.tickets_ext CASCADE; +-- Схема stg для сырого слоя DWH. +CREATE SCHEMA IF NOT EXISTS stg; +-- Внешняя таблица в схеме stg для чтения данных из bookings.tickets через PXF. +DROP EXTERNAL TABLE IF EXISTS stg.tickets_ext; CREATE EXTERNAL TABLE stg.tickets_ext ( ticket_no TEXT, book_ref TEXT, @@ -9,23 +14,19 @@ CREATE EXTERNAL TABLE stg.tickets_ext ( passenger_name TEXT, outbound TEXT ) -LOCATION ('pxf://bookings-db:5432/demo?PROFILE=postgres&SERVER=bookings_db') -FORMAT 'CUSTOM' (FORMATTER='pxfwritable_import') -ENCODING 'UTF8'; +LOCATION ('pxf://bookings.tickets?PROFILE=JDBC&SERVER=bookings-db') +FORMAT 'CUSTOM' (formatter='pxfwritable_import'); --- Создание внутренней таблицы для хранения данных в 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 +-- Внутренняя таблица stg.tickets — сырой слой, все бизнес-колонки как TEXT. +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, + src_created_at_ts TIMESTAMP, + load_dttm TIMESTAMP NOT NULL DEFAULT now(), + batch_id TEXT ) +WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) DISTRIBUTED BY (ticket_no); diff --git a/sql/stg/tickets_dq.sql b/sql/stg/tickets_dq.sql index d0db1af..b5f5feb 100644 --- a/sql/stg/tickets_dq.sql +++ b/sql/stg/tickets_dq.sql @@ -7,17 +7,15 @@ DECLARE v_stg_count BIGINT; BEGIN -- Количество в источнике (новые билеты) + -- Используем только внешние таблицы: stg.tickets_ext + stg.bookings_ext 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; + 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 diff --git a/sql/stg/tickets_load.sql b/sql/stg/tickets_load.sql index 44b9627..73a3ad5 100644 --- a/sql/stg/tickets_load.sql +++ b/sql/stg/tickets_load.sql @@ -21,25 +21,16 @@ SELECT 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 +JOIN stg.bookings_ext AS b ON ext.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' +) AND NOT EXISTS ( -- Защита от дубликатов в рамках одного батча SELECT 1