From d6fed5c7b94fd074075a5225dd3d0af6f28e3d4d Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 17 Jan 2026 21:43:02 +0300 Subject: [PATCH] =?UTF-8?q?=D0=9F=D0=BE=D0=BB=D0=B8=D1=80=D0=BE=D0=B2?= =?UTF-8?q?=D0=BA=D0=B0=20=D0=BA=D0=BE=D0=B4=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- sql/stg/bookings_ddl.sql | 6 ++++++ sql/stg/bookings_dq.sql | 29 +++++++++++++++++++++++++++++ sql/stg/bookings_load.sql | 34 +++++++++++++++++++--------------- sql/stg/tickets_ddl.sql | 7 ++++++- sql/stg/tickets_dq.sql | 37 ++++++++++++++++++++++++++++++++----- sql/stg/tickets_load.sql | 22 +++++++++++++--------- 6 files changed, 105 insertions(+), 30 deletions(-) diff --git a/sql/stg/bookings_ddl.sql b/sql/stg/bookings_ddl.sql index 5d4b8ba..465e6b8 100644 --- a/sql/stg/bookings_ddl.sql +++ b/sql/stg/bookings_ddl.sql @@ -25,5 +25,11 @@ CREATE TABLE IF NOT EXISTS stg.bookings ( batch_id TEXT NOT NULL ) WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +-- Ключ распределения: book_ref +-- Обоснование: book_ref — это уникальный идентификатор бронирования. +-- Использование book_ref обеспечивает: +-- 1. Равномерное распределение данных по сегментам (book_ref имеет высокую кардинальность) +-- 2. Co-location данных bookings и tickets при JOIN по book_ref +-- 3. Оптимизацию запросов, которые фильтруют или группируют по book_ref DISTRIBUTED BY (book_ref); diff --git a/sql/stg/bookings_dq.sql b/sql/stg/bookings_dq.sql index 7d8cc06..24f0e5e 100644 --- a/sql/stg/bookings_dq.sql +++ b/sql/stg/bookings_dq.sql @@ -9,6 +9,8 @@ DECLARE v_prev_ts timestamp; v_src_count bigint; v_stg_count bigint; + v_dup_count bigint; + v_null_amount_count bigint; BEGIN -- Опорная метка: максимум src_created_at_ts среди предыдущих батчей SELECT max(src_created_at_ts) @@ -42,6 +44,33 @@ BEGIN v_stg_count; END IF; + -- Проверка на дубликаты book_ref + SELECT COUNT(*) - COUNT(DISTINCT book_ref) + INTO v_dup_count + FROM stg.bookings AS b + WHERE b.batch_id = v_batch_id; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены дубликаты book_ref (batch_id=%): %', + v_batch_id, + v_dup_count; + END IF; + + -- Проверка на NULL или пустые total_amount + SELECT COUNT(*) + INTO v_null_amount_count + FROM stg.bookings AS b + WHERE b.batch_id = v_batch_id + AND (b.total_amount IS NULL OR b.total_amount = ''); + + IF v_null_amount_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены bookings с NULL или пустым total_amount (batch_id=%): %', + v_batch_id, + v_null_amount_count; + END IF; + RAISE NOTICE 'Проверка количества строк пройдена: источник=%, stg=%', v_src_count, diff --git a/sql/stg/bookings_load.sql b/sql/stg/bookings_load.sql index 55cd551..522c0fc 100644 --- a/sql/stg/bookings_load.sql +++ b/sql/stg/bookings_load.sql @@ -3,6 +3,13 @@ -- берём строки, где book_date больше максимального src_created_at_ts -- среди "старых" батчей; верхняя граница по дате не используется. +-- CTE для определения максимальной даты загрузки предыдущего батча +WITH max_batch_ts AS ( + SELECT COALESCE(MAX(src_created_at_ts), TIMESTAMP '1900-01-01 00:00:00') AS max_ts + FROM stg.bookings + WHERE batch_id <> '{{ run_id }}'::text + OR batch_id IS NULL +) INSERT INTO stg.bookings ( book_ref, book_date, @@ -19,18 +26,15 @@ SELECT now(), '{{ run_id }}'::text FROM stg.bookings_ext AS ext -WHERE ext.book_date > COALESCE( - ( - SELECT max(src_created_at_ts) - FROM stg.bookings - WHERE batch_id <> '{{ run_id }}'::text - OR batch_id IS NULL - ), - TIMESTAMP '1900-01-01 00:00:00' -) - AND NOT EXISTS ( - SELECT 1 - FROM stg.bookings AS b - WHERE b.batch_id = '{{ run_id }}'::text - AND b.book_ref = ext.book_ref::text - ); +CROSS JOIN max_batch_ts AS mb +WHERE ext.book_date > mb.max_ts +AND NOT EXISTS ( + SELECT 1 + FROM stg.bookings AS b + WHERE b.batch_id = '{{ run_id }}'::text + AND b.book_ref = ext.book_ref::text +); + +-- Обновляем статистику для оптимизатора Greenplum +-- Это критично для корректной работы оптимизатора и выбора оптимального плана выполнения +ANALYZE stg.bookings; diff --git a/sql/stg/tickets_ddl.sql b/sql/stg/tickets_ddl.sql index 553cb04..87ea0ef 100644 --- a/sql/stg/tickets_ddl.sql +++ b/sql/stg/tickets_ddl.sql @@ -29,6 +29,11 @@ CREATE TABLE IF NOT EXISTS stg.tickets ( batch_id TEXT ) WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) --- Распределяем по book_ref, чтобы джойны tickets → bookings по book_ref были без motion. +-- Ключ распределения: book_ref +-- Обоснование: book_ref — это основной бизнес-ключ для бронирований. +-- Использование book_ref обеспечивает: +-- 1. Co-location данных tickets и bookings при JOIN по book_ref +-- 2. Равномерное распределение данных по сегментам (book_ref имеет высокую кардинальность) +-- 3. Оптимизацию запросов, которые фильтруют или группируют по book_ref DISTRIBUTED BY (book_ref); diff --git a/sql/stg/tickets_dq.sql b/sql/stg/tickets_dq.sql index 64f63e9..e42a412 100644 --- a/sql/stg/tickets_dq.sql +++ b/sql/stg/tickets_dq.sql @@ -8,6 +8,8 @@ DECLARE v_stg_count BIGINT; v_orphan_count BIGINT; v_null_count BIGINT; + v_dup_count BIGINT; + v_empty_name_count BIGINT; BEGIN -- Опорная метка: максимум src_created_at_ts среди предыдущих батчей SELECT max(src_created_at_ts) @@ -44,15 +46,13 @@ BEGIN END IF; -- Проверка ссылочной целостности: все tickets должны иметь соответствующие bookings в этом же STG + -- Используем LEFT JOIN вместо NOT EXISTS для лучшей производительности на больших объёмах SELECT COUNT(*) INTO v_orphan_count FROM stg.tickets AS t + LEFT JOIN stg.bookings AS b ON t.book_ref = b.book_ref WHERE t.batch_id = v_batch_id - AND NOT EXISTS ( - SELECT 1 - FROM stg.bookings AS b - WHERE b.book_ref = t.book_ref - ); + AND b.book_ref IS NULL; IF v_orphan_count <> 0 THEN RAISE EXCEPTION @@ -75,6 +75,33 @@ BEGIN v_null_count; END IF; + -- Проверка на дубликаты ticket_no + SELECT COUNT(*) - COUNT(DISTINCT ticket_no) + INTO v_dup_count + FROM stg.tickets AS t + WHERE t.batch_id = v_batch_id; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены дубликаты ticket_no (batch_id=%): %', + v_batch_id, + v_dup_count; + END IF; + + -- Проверка на пустые passenger_name + SELECT COUNT(*) + INTO v_empty_name_count + FROM stg.tickets AS t + WHERE t.batch_id = v_batch_id + AND (t.passenger_name IS NULL OR t.passenger_name = ''); + + IF v_empty_name_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены tickets с пустым именем пассажира (batch_id=%): %', + v_batch_id, + v_empty_name_count; + END IF; + RAISE NOTICE 'DQ PASSED: tickets ок (batch_id=%): source=% stg=%', v_batch_id, diff --git a/sql/stg/tickets_load.sql b/sql/stg/tickets_load.sql index 3a32059..d90027d 100644 --- a/sql/stg/tickets_load.sql +++ b/sql/stg/tickets_load.sql @@ -1,6 +1,13 @@ -- Загрузка инкремента из stg.tickets_ext в stg.tickets -- Инкремент определяется по дате бронирования (book_date из bookings.bookings) +-- CTE для определения максимальной даты загрузки предыдущего батча +WITH max_batch_ts AS ( + SELECT COALESCE(MAX(src_created_at_ts), TIMESTAMP '1900-01-01 00:00:00') AS max_ts + FROM stg.tickets + WHERE batch_id <> '{{ run_id }}'::text + OR batch_id IS NULL +) INSERT INTO stg.tickets ( ticket_no, book_ref, @@ -22,18 +29,15 @@ SELECT '{{ run_id }}'::text FROM stg.tickets_ext AS ext 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' -) +CROSS JOIN max_batch_ts AS mb +WHERE b.book_date > mb.max_ts AND NOT EXISTS ( -- Защита от дублей: ticket_no в источнике уникален, и в stg его не дублируем. SELECT 1 FROM stg.tickets AS t WHERE t.ticket_no = ext.ticket_no ); + +-- Обновляем статистику для оптимизатора Greenplum +-- Это критично для корректной работы оптимизатора и выбора оптимального плана выполнения +ANALYZE stg.tickets;