Проверка/рецензирование доработки
This commit is contained in:
@@ -80,6 +80,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;
|
||||||
```
|
```
|
||||||
|
|
||||||
@@ -89,9 +90,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):
|
||||||
|
|
||||||
|
|||||||
@@ -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",
|
||||||
|
|||||||
@@ -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"],
|
||||||
|
|||||||
@@ -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,9 +27,7 @@ make up
|
|||||||
make bookings-init
|
make bookings-init
|
||||||
```
|
```
|
||||||
|
|
||||||
Важно: генератор `bookings` в этом стенде поддерживается только в режиме `BOOKINGS_JOBS=1`.
|
3) В Greenplum созданы STG‑объекты `stg.bookings_ext`/`stg.bookings` и `stg.tickets_ext`/`stg.tickets` (выберите один вариант):
|
||||||
|
|
||||||
3) В Greenplum созданы `stg.bookings_ext` и `stg.bookings` (выберите один вариант):
|
|
||||||
|
|
||||||
- учебный вариант: запустить DAG `bookings_stg_ddl` в Airflow UI;
|
- учебный вариант: запустить DAG `bookings_stg_ddl` в Airflow UI;
|
||||||
- технический шорткат: `make ddl-gp`.
|
- технический шорткат: `make ddl-gp`.
|
||||||
@@ -70,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`
|
||||||
|
|
||||||
- логирует краткую сводку в конце запуска.
|
- логирует краткую сводку в конце запуска.
|
||||||
|
|
||||||
|
|||||||
+79
-50
@@ -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. Результаты тестирования
|
||||||
|
|
||||||
|
|||||||
@@ -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
@@ -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 $$;
|
||||||
|
|||||||
@@ -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
|
|
||||||
);
|
);
|
||||||
|
|||||||
Reference in New Issue
Block a user