Первая версия dag

This commit is contained in:
2026-01-17 15:09:08 +03:00
parent 0208e6fd71
commit 464e631eea
4 changed files with 152 additions and 3 deletions
+18 -3
View File
@@ -54,7 +54,7 @@ with DAG(
template_searchpath="/sql", template_searchpath="/sql",
default_args=default_args, default_args=default_args,
tags=["demo", "bookings", "greenplum", "stg"], tags=["demo", "bookings", "greenplum", "stg"],
description="Учебный DAG: загрузка из bookings-db в stg.bookings (Greenplum)", description="Учебный DAG: загрузка из bookings-db в stg.bookings и stg.tickets (Greenplum)",
) as dag: ) as dag:
# 1. Генерируем один (или несколько стартовых) учебный день в демо-БД bookings # 1. Генерируем один (или несколько стартовых) учебный день в демо-БД bookings
generate_bookings_day = PostgresOperator( generate_bookings_day = PostgresOperator(
@@ -78,10 +78,25 @@ with DAG(
sql="stg/bookings_dq.sql", 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( finish_summary = PythonOperator(
task_id="finish_summary", task_id="finish_summary",
python_callable=_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
+31
View File
@@ -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);
+54
View File
@@ -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
+49
View File
@@ -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
);