Проверка/рецензирование доработки

This commit is contained in:
2026-01-17 19:36:02 +03:00
parent 93ba59c1c9
commit f36a98497a
8 changed files with 189 additions and 109 deletions
+4 -3
View File
@@ -78,6 +78,7 @@ make bookings-init
make gp-psql make gp-psql
-- внутри psql: -- внутри psql:
SELECT COUNT(*) FROM stg.bookings; SELECT COUNT(*) FROM stg.bookings;
SELECT COUNT(*) FROM stg.tickets;
SELECT * FROM stg.bookings ORDER BY src_created_at_ts DESC LIMIT 10; SELECT * FROM stg.bookings ORDER BY src_created_at_ts DESC LIMIT 10;
``` ```
@@ -87,9 +88,9 @@ SELECT * FROM stg.bookings ORDER BY src_created_at_ts DESC LIMIT 10;
Основные (для потока bookings → DWH): Основные (для потока bookings → DWH):
- `bookings_stg_ddl` — создаёт `stg.bookings_ext` и `stg.bookings` в Greenplum; - `bookings_stg_ddl` — создаёт `stg.bookings_ext`/`stg.bookings` и `stg.tickets_ext`/`stg.tickets` в Greenplum;
- `bookings_to_gp_stage` — генерирует учебный день в `bookings-db` и грузит инкремент в `stg.bookings` - `bookings_to_gp_stage` — генерирует учебный день в `bookings-db`, грузит инкремент в `stg.bookings` и `stg.tickets`
(через PXF), затем выполняет DQ‑проверку. (через PXF), затем выполняет DQ‑проверки.
Вспомогательные (побочный трек с CSV): Вспомогательные (побочный трек с CSV):
+3 -2
View File
@@ -5,8 +5,9 @@ from __future__ import annotations
Запускается вручную перед DAG загрузки bookings_to_gp_stage или после изменения DDL. Запускается вручную перед DAG загрузки bookings_to_gp_stage или после изменения DDL.
""" """
from datetime import datetime, timedelta from datetime import timedelta
import pendulum
from airflow.providers.postgres.operators.postgres import PostgresOperator from airflow.providers.postgres.operators.postgres import PostgresOperator
from airflow import DAG from airflow import DAG
@@ -17,7 +18,7 @@ default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(secon
with DAG( with DAG(
dag_id="bookings_stg_ddl", dag_id="bookings_stg_ddl",
start_date=datetime(2024, 1, 1), start_date=pendulum.datetime(2024, 1, 1, tz="UTC"),
schedule=None, schedule=None,
catchup=False, catchup=False,
template_searchpath="/sql", template_searchpath="/sql",
+4 -2
View File
@@ -15,8 +15,9 @@ from __future__ import annotations
- `run_id` используется как метка запуска (в `batch_id`, в логах и DQ). - `run_id` используется как метка запуска (в `batch_id`, в логах и DQ).
""" """
from datetime import datetime, timedelta from datetime import timedelta
import pendulum
from airflow.operators.python import PythonOperator from airflow.operators.python import PythonOperator
from airflow.providers.postgres.operators.postgres import PostgresOperator from airflow.providers.postgres.operators.postgres import PostgresOperator
@@ -48,9 +49,10 @@ def _finish_summary() -> None:
with DAG( with DAG(
dag_id="bookings_to_gp_stage", dag_id="bookings_to_gp_stage",
start_date=datetime(2024, 1, 1), start_date=pendulum.datetime(2024, 1, 1, tz="UTC"),
schedule=None, schedule=None,
catchup=False, catchup=False,
max_active_runs=1,
template_searchpath="/sql", template_searchpath="/sql",
default_args=default_args, default_args=default_args,
tags=["demo", "bookings", "greenplum", "stg"], tags=["demo", "bookings", "greenplum", "stg"],
+18 -3
View File
@@ -2,7 +2,7 @@
Этот DAG — основной учебный пример в стенде. Он показывает путь данных из источника **Postgres** Этот DAG — основной учебный пример в стенде. Он показывает путь данных из источника **Postgres**
(`bookings-db`, демо‑БД `demo`) в сырой слой **STG** в **Greenplum** с инкрементальной загрузкой (`bookings-db`, демо‑БД `demo`) в сырой слой **STG** в **Greenplum** с инкрементальной загрузкой
и простой проверкой качества данных. и простыми проверками качества данных.
## Что делает DAG ## Что делает DAG
@@ -10,6 +10,8 @@
(генератор всегда “шагает” вперёд от `max(book_date)`). (генератор всегда “шагает” вперёд от `max(book_date)`).
- В Greenplum загружает инкремент в `stg.bookings` через внешнюю таблицу `stg.bookings_ext`, используя PXF. - В Greenplum загружает инкремент в `stg.bookings` через внешнюю таблицу `stg.bookings_ext`, используя PXF.
- Сверяет количество строк между источником (за окно инкремента) и загруженным батчем. - Сверяет количество строк между источником (за окно инкремента) и загруженным батчем.
- Загружает инкремент в `stg.tickets` через внешнюю таблицу `stg.tickets_ext`, используя PXF.
- Запускает DQ‑проверки для `stg.tickets` (количество, ссылочная целостность, обязательные поля).
## Что должно быть готово перед запуском ## Что должно быть готово перед запуском
@@ -25,7 +27,7 @@ make up
make bookings-init make bookings-init
``` ```
3) В Greenplum созданы `stg.bookings_ext` и `stg.bookings` (выберите один вариант): 3) В Greenplum созданы STG‑объекты `stg.bookings_ext`/`stg.bookings` и `stg.tickets_ext`/`stg.tickets` (выберите один вариант):
- учебный вариант: запустить DAG `bookings_stg_ddl` в Airflow UI; - учебный вариант: запустить DAG `bookings_stg_ddl` в Airflow UI;
- технический шорткат: `make ddl-gp`. - технический шорткат: `make ddl-gp`.
@@ -68,7 +70,20 @@ make bookings-init
вставленных в `stg.bookings` для текущего `batch_id`; вставленных в `stg.bookings` для текущего `batch_id`;
- при расхождении делает `RAISE EXCEPTION` с понятным текстом. - при расхождении делает `RAISE EXCEPTION` с понятным текстом.
4) `finish_summary` 4) `load_tickets_to_stg`
- выполняет `sql/stg/tickets_load.sql` в Greenplum;
- так как в `bookings.tickets` нет явной временной колонки, окно инкремента берётся по `book_date`
из связанной внешней таблицы `stg.bookings_ext` (JOIN по `book_ref`);
- вставляет строки в `stg.tickets`, добавляя `src_created_at_ts`, `load_dttm` и `batch_id={{ run_id }}`.
5) `check_tickets_dq`
- выполняет `sql/stg/tickets_dq.sql` в Greenplum;
- проверяет количество строк в том же окне инкремента, а также ссылочную целостность и обязательные поля;
- при проблемах делает `RAISE EXCEPTION`, чтобы DAG падал “красным”.
6) `finish_summary`
- логирует краткую сводку в конце запуска. - логирует краткую сводку в конце запуска.
+80 -51
View File
@@ -37,27 +37,26 @@
```sql ```sql
-- sql/stg/tickets_ddl.sql -- sql/stg/tickets_ddl.sql
DROP TABLE IF EXISTS stg.tickets CASCADE; CREATE TABLE IF NOT EXISTS stg.tickets (
CREATE TABLE stg.tickets (
-- Бизнес-атрибуты (из источника) -- Бизнес-атрибуты (из источника)
ticket_no TEXT NOT NULL, -- номер билета ticket_no TEXT NOT NULL, -- номер билета
book_ref TEXT NOT NULL, -- ссылка на бронирование book_ref TEXT NOT NULL, -- ссылка на бронирование
passenger_id TEXT, -- идентификатор пассажира passenger_id TEXT, -- идентификатор пассажира
passenger_name TEXT, -- имя пассажира passenger_name TEXT, -- имя пассажира
outbound TEXT, -- флаг направления (boolean → text для Greenplum) outbound TEXT, -- флаг направления (в источнике boolean; в STG храним как TEXT)
-- Технические атрибуты -- Технические атрибуты
src_created_at_ts TIMESTAMP, -- временная метка источника (для инкремента) src_created_at_ts TIMESTAMP, -- временная метка источника (для инкремента)
load_dttm TIMESTAMP NOT NULL, -- время загрузки в Greenplum load_dttm TIMESTAMP NOT NULL DEFAULT now(), -- время загрузки в Greenplum
batch_id TEXT -- идентификатор батча (из Airflow run_id) batch_id TEXT -- идентификатор батча (из Airflow run_id)
) )
DISTRIBUTED BY (ticket_no); -- распределение по первичному ключу WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
DISTRIBUTED BY (book_ref); -- распределение по ключу связи с bookings
``` ```
**Решения по проектированию:** **Решения по проектированию:**
- **DISTRIBUTED BY (ticket_no):** равномерное распределение, так как это первичный ключ - **DISTRIBUTED BY (book_ref):** типовой джойн `tickets → bookings` идёт по `book_ref`, так меньше motion в MPP
- **boolean → TEXT:** Greenplum хорошо работает с текстовым представлением для флагов - **boolean → TEXT:** в сыром STG храним бизнес-колонки как TEXT (для обучения и минимизации кастов на входе)
- **src_created_at_ts:** временная метка из даты связанного бронирования (см. раздел 7) - **src_created_at_ts:** временная метка из даты связанного бронирования (см. раздел 7)
## 4. Проектирование внешней таблицы (PXF) ## 4. Проектирование внешней таблицы (PXF)
@@ -122,7 +121,11 @@ CREATE TABLE IF NOT EXISTS stg.tickets (
batch_id TEXT batch_id TEXT
) )
WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
DISTRIBUTED BY (ticket_no); -- Распределяем по book_ref, чтобы джойны tickets → bookings по book_ref были без motion.
DISTRIBUTED BY (book_ref);
-- На случай, если таблица уже была создана раньше с другим ключом распределения.
ALTER TABLE IF EXISTS stg.tickets SET DISTRIBUTED BY (book_ref);
``` ```
## 6. LOAD-скрипт (загрузка инкремента) ## 6. LOAD-скрипт (загрузка инкремента)
@@ -173,11 +176,10 @@ WHERE b.book_date > COALESCE(
TIMESTAMP '1900-01-01 00:00:00' TIMESTAMP '1900-01-01 00:00:00'
) )
AND NOT EXISTS ( AND NOT EXISTS (
-- Защита от дубликатов в рамках одного батча -- Защита от дублей: ticket_no в источнике уникален, и в stg его не дублируем.
SELECT 1 SELECT 1
FROM stg.tickets AS t FROM stg.tickets AS t
WHERE t.batch_id = '{{ run_id }}'::text WHERE t.ticket_no = ext.ticket_no
AND t.ticket_no = ext.ticket_no
); );
``` ```
@@ -194,56 +196,83 @@ AND NOT EXISTS (
Скрипт `sql/stg/tickets_dq.sql`: Скрипт `sql/stg/tickets_dq.sql`:
```sql ```sql
-- Проверка 1: совпадение количества билетов в источнике и STG -- Проверки качества данных для tickets
-- Используем только внешние таблицы: stg.tickets_ext + stg.bookings_ext
DO $$ DO $$
DECLARE DECLARE
v_batch_id TEXT := '{{ run_id }}'::text;
v_prev_ts TIMESTAMP;
v_source_count BIGINT; v_source_count BIGINT;
v_stg_count BIGINT; v_stg_count BIGINT;
v_orphan_count BIGINT;
v_null_count BIGINT;
BEGIN BEGIN
-- Количество в источнике (новые билеты) -- Опорная метка: максимум src_created_at_ts среди предыдущих батчей
SELECT COUNT(*) INTO v_source_count SELECT max(src_created_at_ts)
INTO v_prev_ts
FROM stg.tickets
WHERE batch_id <> v_batch_id
OR batch_id IS NULL;
-- Источник: считаем строки в том же окне инкремента, что и загрузка
SELECT COUNT(*)
INTO v_source_count
FROM stg.tickets_ext AS t FROM stg.tickets_ext AS t
JOIN stg.bookings_ext AS b ON t.book_ref = b.book_ref JOIN stg.bookings_ext AS b ON t.book_ref = b.book_ref
WHERE b.book_date > COALESCE( WHERE b.book_date > COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00');
(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 (текущий батч) IF v_source_count = 0 THEN
SELECT COUNT(*) INTO v_stg_count RAISE EXCEPTION
FROM stg.tickets 'В источнике tickets_ext нет строк для окна инкремента (book_date > %). Проверьте генерацию данных (таск generate_bookings_day).',
WHERE batch_id = '{{ run_id }}'::text; COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00');
-- Проверка совпадения
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 IF;
END $$;
-- Проверка 2: ссылочная целостность (все ticket_no должны иметь соответствующие book_ref) -- STG: считаем строки текущего батча
SELECT SELECT COUNT(*)
COUNT(*) AS orphan_tickets INTO v_stg_count
FROM stg.tickets FROM stg.tickets
WHERE batch_id = '{{ run_id }}'::text WHERE batch_id = v_batch_id;
IF v_source_count <> v_stg_count THEN
RAISE EXCEPTION
'DQ FAILED: несовпадение количества билетов. Источник: %, STG (batch_id=%): %',
v_source_count,
v_batch_id,
v_stg_count;
END IF;
-- Ссылочная целостность: tickets должны иметь соответствующие bookings в STG
SELECT COUNT(*)
INTO v_orphan_count
FROM stg.tickets AS t
WHERE t.batch_id = v_batch_id
AND NOT EXISTS ( AND NOT EXISTS (
SELECT 1 SELECT 1
FROM stg.bookings FROM stg.bookings AS b
WHERE stg.bookings.book_ref = stg.tickets.book_ref WHERE b.book_ref = t.book_ref
); );
-- Ожидаемое значение: 0
-- Проверка 3: отсутствие NULL в обязательных полях IF v_orphan_count <> 0 THEN
SELECT RAISE EXCEPTION
COUNT(*) AS null_tickets 'DQ FAILED: найдены tickets без соответствующих bookings (batch_id=%): %',
FROM stg.tickets v_batch_id,
WHERE batch_id = '{{ run_id }}'::text v_orphan_count;
AND (ticket_no IS NULL OR book_ref IS NULL); END IF;
-- Ожидаемое значение: 0
-- Обязательные поля
SELECT COUNT(*)
INTO v_null_count
FROM stg.tickets AS t
WHERE t.batch_id = v_batch_id
AND (t.ticket_no IS NULL OR t.book_ref IS NULL);
IF v_null_count <> 0 THEN
RAISE EXCEPTION
'DQ FAILED: найдены tickets с NULL в обязательных полях (batch_id=%): %',
v_batch_id,
v_null_count;
END IF;
END $$;
``` ```
## 8. Интеграция с существующими DAG ## 8. Интеграция с существующими DAG
@@ -396,7 +425,7 @@ WHERE b.book_ref IS NULL;
8. ✅ Проверка через Airflow UI 8. ✅ Проверка через Airflow UI
9. ✅ Проверка идемпотентности DDL 9. ✅ Проверка идемпотентности DDL
10. ✅ Проверка инкрементальной загрузки 10. ✅ Проверка инкрементальной загрузки
11. Обновить `README.md` (добавить tickets в список STG-таблиц) 11. Обновить `README.md` (добавить tickets в список STG-таблиц)
## 11. Сопутствующие изменения ## 11. Сопутствующие изменения
@@ -412,8 +441,8 @@ WHERE b.book_ref IS NULL;
### Необходимые изменения в документации: ### Необходимые изменения в документации:
**README.md:** **README.md:**
- Добавить `tickets` в список STG-таблиц - Добавить `tickets` в список STG-таблиц
- Обновить описание DAG `bookings_to_gp_stage` (упомянуть загрузку билетов) - Обновить описание DAG `bookings_to_gp_stage` (упомянуть загрузку билетов)
## 12. Результаты тестирования ## 12. Результаты тестирования
+3 -1
View File
@@ -29,4 +29,6 @@ CREATE TABLE IF NOT EXISTS stg.tickets (
batch_id TEXT batch_id TEXT
) )
WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
DISTRIBUTED BY (ticket_no); -- Распределяем по book_ref, чтобы джойны tickets → bookings по book_ref были без motion.
DISTRIBUTED BY (book_ref);
+63 -32
View File
@@ -1,52 +1,83 @@
-- Проверки качества данных для tickets -- Проверки качества данных для tickets
-- Проверка 1: совпадение количества билетов в источнике и STG
DO $$ DO $$
DECLARE DECLARE
v_batch_id TEXT := '{{ run_id }}'::text;
v_prev_ts TIMESTAMP;
v_source_count BIGINT; v_source_count BIGINT;
v_stg_count BIGINT; v_stg_count BIGINT;
v_orphan_count BIGINT;
v_null_count BIGINT;
BEGIN BEGIN
-- Количество в источнике (новые билеты) -- Опорная метка: максимум src_created_at_ts среди предыдущих батчей
-- Используем только внешние таблицы: stg.tickets_ext + stg.bookings_ext SELECT max(src_created_at_ts)
SELECT COUNT(*) INTO v_source_count INTO v_prev_ts
FROM stg.tickets
WHERE batch_id <> v_batch_id
OR batch_id IS NULL;
-- Количество в источнике (новые билеты в том же окне инкремента, что и загрузка)
SELECT COUNT(*)
INTO v_source_count
FROM stg.tickets_ext AS t FROM stg.tickets_ext AS t
JOIN stg.bookings_ext AS b ON t.book_ref = b.book_ref JOIN stg.bookings_ext AS b ON t.book_ref = b.book_ref
WHERE b.book_date > COALESCE( WHERE b.book_date > COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00');
(SELECT max(src_created_at_ts) FROM stg.tickets
WHERE batch_id <> '{{ run_id }}'::text OR batch_id IS NULL), IF v_source_count = 0 THEN
TIMESTAMP '1900-01-01 00:00:00' RAISE EXCEPTION
); 'В источнике tickets_ext нет строк для окна инкремента (book_date > %). Проверьте генерацию данных (таск generate_bookings_day).',
COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00');
END IF;
-- Количество в STG (текущий батч) -- Количество в STG (текущий батч)
SELECT COUNT(*) INTO v_stg_count SELECT COUNT(*)
INTO v_stg_count
FROM stg.tickets FROM stg.tickets
WHERE batch_id = '{{ run_id }}'::text; WHERE batch_id = v_batch_id;
-- Проверка совпадения
IF v_source_count <> v_stg_count THEN IF v_source_count <> v_stg_count THEN
RAISE EXCEPTION 'DQ FAILED: несовпадение количества билетов. Источник: %, STG: %', RAISE EXCEPTION
v_source_count, v_stg_count; 'DQ FAILED: несовпадение количества билетов. Источник: %, STG (batch_id=%): %',
ELSE v_source_count,
RAISE NOTICE 'DQ PASSED: количество билетов совпадает (%)', v_stg_count; v_batch_id,
v_stg_count;
END IF; END IF;
END $$;
-- Проверка 2: ссылочная целостность (все ticket_no должны иметь соответствующие book_ref) -- Проверка ссылочной целостности: все tickets должны иметь соответствующие bookings в этом же STG
SELECT SELECT COUNT(*)
COUNT(*) AS orphan_tickets INTO v_orphan_count
FROM stg.tickets FROM stg.tickets AS t
WHERE batch_id = '{{ run_id }}'::text WHERE t.batch_id = v_batch_id
AND NOT EXISTS ( AND NOT EXISTS (
SELECT 1 SELECT 1
FROM stg.bookings FROM stg.bookings AS b
WHERE stg.bookings.book_ref = stg.tickets.book_ref WHERE b.book_ref = t.book_ref
); );
-- Ожидаемое значение: 0
-- Проверка 3: отсутствие NULL в обязательных полях IF v_orphan_count <> 0 THEN
SELECT RAISE EXCEPTION
COUNT(*) AS null_tickets 'DQ FAILED: найдены tickets без соответствующих bookings (batch_id=%): %',
FROM stg.tickets v_batch_id,
WHERE batch_id = '{{ run_id }}'::text v_orphan_count;
AND (ticket_no IS NULL OR book_ref IS NULL); END IF;
-- Ожидаемое значение: 0
-- Проверка обязательных полей
SELECT COUNT(*)
INTO v_null_count
FROM stg.tickets AS t
WHERE t.batch_id = v_batch_id
AND (t.ticket_no IS NULL OR t.book_ref IS NULL);
IF v_null_count <> 0 THEN
RAISE EXCEPTION
'DQ FAILED: найдены tickets с NULL в обязательных полях (batch_id=%): %',
v_batch_id,
v_null_count;
END IF;
RAISE NOTICE
'DQ PASSED: tickets ок (batch_id=%): source=% stg=%',
v_batch_id,
v_source_count,
v_stg_count;
END $$;
+2 -3
View File
@@ -32,9 +32,8 @@ WHERE b.book_date > COALESCE(
TIMESTAMP '1900-01-01 00:00:00' TIMESTAMP '1900-01-01 00:00:00'
) )
AND NOT EXISTS ( AND NOT EXISTS (
-- Защита от дубликатов в рамках одного батча -- Защита от дублей: ticket_no в источнике уникален, и в stg его не дублируем.
SELECT 1 SELECT 1
FROM stg.tickets AS t FROM stg.tickets AS t
WHERE t.batch_id = '{{ run_id }}'::text WHERE t.ticket_no = ext.ticket_no
AND t.ticket_no = ext.ticket_no
); );