Первая версия dag
This commit is contained in:
@@ -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
|
||||||
|
|||||||
@@ -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);
|
||||||
@@ -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
|
||||||
@@ -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
|
||||||
|
);
|
||||||
Reference in New Issue
Block a user