From cacf989a9cbf39fb1c4a718612c4a53a3e8331ed Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Wed, 10 Dec 2025 16:59:50 +0300 Subject: [PATCH] =?UTF-8?q?=D0=B2=D1=8B=D0=BD=D0=B5=D1=81=D0=BB=D0=B8=20sq?= =?UTF-8?q?l=20=D0=B8=D0=B7=20dag=20=D0=B2=20=D0=BF=D0=B0=D0=BF=D0=BA?= =?UTF-8?q?=D1=83=20sql?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- AGENTS.md | 16 + README.md | 24 ++ airflow/dags/bookings_to_gp_stage.py | 393 ++----------------- docker-compose.yml | 3 + docs/internal/bookings_stg_design.md | 52 +-- docs/internal/bookings_stg_readme.md | 31 +- sql/ddl_gp.sql | 27 +- sql/src/bookings_generate_day_if_missing.sql | 59 +++ sql/stg/bookings_ddl.sql | 29 ++ sql/stg/bookings_dq.sql | 50 +++ sql/stg/bookings_load.sql | 32 ++ 11 files changed, 292 insertions(+), 424 deletions(-) create mode 100644 sql/src/bookings_generate_day_if_missing.sql create mode 100644 sql/stg/bookings_ddl.sql create mode 100644 sql/stg/bookings_dq.sql create mode 100644 sql/stg/bookings_load.sql diff --git a/AGENTS.md b/AGENTS.md index bbfc4ce..1cb3e3e 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -34,6 +34,22 @@ - Форматирование: `black` (88 cols) и `isort`. Если не уверены — запустите `make fmt`. - Язык: комментарии, docstring и документацию — на русском; имена идентификаторов — на английском. +### Структура и нейминг SQL (слои DWH) +- В каталоге `sql/` придерживаемся слоёв DWH: + - `sql/src/` — скрипты, работающие с исходными системами (например, `bookings_generate_day_if_missing.sql`); + - `sql/stg/` — скрипты для стейджинга (`bookings_ddl.sql`, `bookings_load.sql`, `bookings_dq.sql`); + - в будущем можно добавить `sql/ods/`, `sql/dds/`, `sql/dm/` по мере роста стенда. +- Именование файлов: `{объект}_{роль}.sql`, где: + - `объект` — логическое имя сущности (`bookings`, `orders`, и т.п.); + - `роль` — `ddl` (создание/изменение объектов), `load` (загрузка/инкремент), `dq` (проверки качества данных) и т.п. +- Общие DDL-скрипты (например, `sql/ddl_gp.sql`) могут подключать файловые DDL через `\i`, но сами определения таблиц живут рядом с объектом (`sql/stg/bookings_ddl.sql` и т.п.). + +### Airflow + SQL +- В учебных DAG’ах, где основная логика — в SQL, по умолчанию используем `PostgresOperator` + Airflow Connections: + - DAG оркестрирует шаги и подключение к БД; + - SQL-скрипты лежат в `sql/...` и подключаются по пути (`sql='sql/stg/bookings_load.sql'`). +- Сложную ручную работу с подключениями (`psycopg2`, ENV-фоллбеки) используем только там, где реально много Python-логики и это помогает учебной цели. + ## Тестирование - Тесты лежат в `tests/` (pytest). Запуск: `make test`. - Есть юнит‑тесты для `helpers/greenplum.py` и smoke‑тесты DAG‑структуры (`tests/test_dags_smoke.py`). diff --git a/README.md b/README.md index ed19919..ccbe8ac 100644 --- a/README.md +++ b/README.md @@ -239,6 +239,30 @@ docker compose -f docker-compose.yml exec bookings-db bash -lc 'PGPASSWORD="$POS - Объем загруженных данных - Отсутствие дубликатов записей +### Пример DAG с SQL-скриптами (bookings → stg) + +В репозитории есть учебный DAG `bookings_to_gp_stage`, который показывает «канонический» способ работы с SQL в Airflow: + +- подключение к БД через Airflow Connections (`bookings_db`, `greenplum_conn`); +- бизнес-логика инкрементальной загрузки и DQ вынесена в SQL-файлы в каталоге `sql/`: + - `sql/src/bookings_generate_day_if_missing.sql` — генерация учебного дня в демо-БД bookings (идемпотентно); + - `sql/stg/bookings_ddl.sql` — DDL для схемы `stg` и таблиц `stg.bookings_ext` / `stg.bookings`; + - `sql/stg/bookings_load.sql` — загрузка инкремента из `stg.bookings_ext` в `stg.bookings`; + - `sql/stg/bookings_dq.sql` — проверка количества строк между источником и stg. + +Фрагмент DAG: + +```python +load_bookings_to_stg = PostgresOperator( + task_id="load_bookings_to_stg", + postgres_conn_id="greenplum_conn", + sql="/sql/stg/bookings_load.sql", + params={"load_date": "{{ ds }}", "batch_id": "{{ ds_nodash }}"}, +) +``` + +Такой подход помогает держать оркестрацию (DAG) и SQL-логику в отдельных файлах и легче сравнивать её с теорией из статьи про моделирование DWH. + ### Ограничения учебного стенда - **Greenplum** запущен в single-node режиме (для обучения) diff --git a/airflow/dags/bookings_to_gp_stage.py b/airflow/dags/bookings_to_gp_stage.py index 88ad88c..4724671 100644 --- a/airflow/dags/bookings_to_gp_stage.py +++ b/airflow/dags/bookings_to_gp_stage.py @@ -4,19 +4,17 @@ from __future__ import annotations Учебный DAG для менти: показывает, как устроен поток от источника bookings-db до слоя stg в Greenplum. -Важно: этот файл специально должен быть хорошо задокументирован — -docstring и комментарии помогают студенту понять, «зачем» каждая задача, -а не только «что именно она делает». +В этом примере мы сознательно используем: +- Airflow Connections для подключения к БД; +- PostgresOperator, который берёт SQL-скрипты с диска; +чтобы студент увидел «канонический» способ работы с SQL в DAG. """ -import logging -import os from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator - -from helpers.greenplum import get_gp_conn +from airflow.providers.postgres.operators.postgres import PostgresOperator default_args = { @@ -26,341 +24,21 @@ default_args = { } -BOOKINGS_CONN_ID = os.getenv("BOOKINGS_CONN_ID", "bookings_db") - - -def _get_bookings_conn(): - """ - Возвращает подключение к демо-БД bookings. - - Приоритет: - 1. Airflow Connection с ID из BOOKINGS_CONN_ID (по умолчанию: bookings_db) - 2. Прямое подключение по переменным окружения (фоллбек) - """ - try: - from airflow.providers.postgres.hooks.postgres import PostgresHook - - hook = PostgresHook(postgres_conn_id=BOOKINGS_CONN_ID) - conn = hook.get_conn() - logging.info( - "✅ Подключение к bookings-db через Airflow Connection '%s' успешно", - BOOKINGS_CONN_ID, - ) - return conn - except Exception as exc: # pragma: no cover - фоллбек для нестандартных окружений - logging.warning( - "⚠️ Не удалось подключиться к bookings-db через Airflow Connection '%s': %s", - BOOKINGS_CONN_ID, - exc, - ) - logging.info("🔄 Пробуем прямое подключение по переменным окружения") - - import psycopg2 - - conn_params = { - # Внутри Docker-сети bookings-db доступен по имени сервиса и порту 5432 - "host": os.getenv("BOOKINGS_DB_HOST", "bookings-db"), - "port": int(os.getenv("BOOKINGS_DB_PORT_INTERNAL", "5432")), - "dbname": os.getenv("BOOKINGS_DB_NAME", "demo"), - "user": os.getenv("BOOKINGS_DB_USER", "bookings"), - "password": os.getenv("BOOKINGS_DB_PASSWORD", "bookings"), - } - logging.info( - "🔗 Подключение к bookings-db по ENV: %s:%s/%s", - conn_params["host"], - conn_params["port"], - conn_params["dbname"], - ) - return psycopg2.connect(**conn_params) - - -def _generate_bookings_day(load_date: str) -> None: - """ - Готовит данные за указанный день в bookings-db. - - Идея: - - если за load_date уже есть строки в bookings.bookings → ничего не делаем; - - если нет → запускаем генератор (аналог make bookings-generate-day). - """ - logging.info("Запускаем подготовку данных в bookings-db за дату %s", load_date) - - with _get_bookings_conn() as conn, conn.cursor() as cur: - # 1. Проверяем, что демо-БД установлена (таблица bookings.bookings существует) - logging.info("Проверяем наличие таблицы bookings.bookings...") - cur.execute("SELECT to_regclass('bookings.bookings')") - table_regclass = cur.fetchone()[0] - if table_regclass is None: - raise ValueError( - "❌ Таблица bookings.bookings не найдена. " - "Сначала выполните make bookings-init, чтобы подготовить демо-БД." - ) - - # 2. Проверяем, есть ли уже данные за нужный день - logging.info( - "Проверяем, есть ли данные за %s в bookings.bookings...", load_date - ) - cur.execute( - """ - SELECT EXISTS ( - SELECT 1 - FROM bookings.bookings - WHERE book_date::date = %s::date - ) - """, - (load_date,), - ) - has_day = bool(cur.fetchone()[0]) - - if has_day: - logging.info( - "Данные за %s уже есть в bookings.bookings — " - "генерация не требуется (идемпотентность).", - load_date, - ) - return - - logging.info( - "Данных за %s нет — запускаем генератор демобазы " - "(аналог make bookings-generate-day)...", - load_date, - ) - - # 3. Запускаем генерацию следующего дня через тот же DO-блок, - # который используется в скрипте bookings/generate_next_day.sql. - # Это гарантирует, что логика совпадает с CLI-сценарием. - cur.execute( - """ - DO $$ - DECLARE - v_max_book_date timestamptz; - v_start_date timestamptz; - v_end_date timestamptz; - v_jobs integer := COALESCE(current_setting('bookings.jobs', true), '1')::integer; - v_init_days integer := COALESCE(current_setting('bookings.init_days', true), '1')::integer; - v_start_cfg text := COALESCE(current_setting('bookings.start_date', true), '2017-01-01'); - BEGIN - -- Проверяем, что демобаза установлена - IF to_regclass('bookings.bookings') IS NULL THEN - RAISE EXCEPTION 'Таблица bookings.bookings не найдена. Сначала выполните make bookings-init.'; - END IF; - - -- Ищем последнюю сгенерированную дату - SELECT max(book_date) INTO v_max_book_date FROM bookings.bookings; - - IF v_max_book_date IS NULL THEN - -- База пустая: берём стартовую дату из конфигурации (или дефолтную) - v_start_date := date_trunc('day', v_start_cfg::timestamptz); - ELSE - -- Продолжаем с дня, следующего за максимальной датой - v_start_date := date_trunc('day', v_max_book_date) + interval '1 day'; - END IF; - - -- Первая генерация вызывает generate(), последующие — continue() - IF v_max_book_date IS NULL THEN - v_end_date := v_start_date + (v_init_days || ' days')::interval; - CALL generate(v_start_date, v_end_date, v_jobs); - ELSE - v_end_date := v_start_date + interval '1 day'; - CALL continue(v_end_date, v_jobs); - END IF; - - -- Ждём завершения фоновых джобов генератора, чтобы данные успели записаться - WHILE busy() LOOP - PERFORM pg_sleep(1); - END LOOP; - PERFORM dblink_disconnect(unnest(dblink_get_connections())); - END $$; - """ - ) - conn.commit() - - logging.info("Генерация данных за %s в bookings-db завершена.", load_date) - - -def _get_last_loaded_ts_from_gp() -> str | None: - """ - Возвращает максимальное значение src_created_at_ts из stg.bookings. - - Если данных ещё нет, возвращает None — это будет означать - режим полной загрузки (full). - """ - with get_gp_conn() as conn, conn.cursor() as cur: - logging.info("Проверяем наличие таблицы stg.bookings в Greenplum...") - cur.execute("SELECT to_regclass('stg.bookings')") - table_regclass = cur.fetchone()[0] - if table_regclass is None: - raise ValueError( - "❌ Таблица stg.bookings не найдена. " - "Убедитесь, что выполнен DDL для схемы stg (например, make ddl-gp)." - ) - - logging.info( - "Читаем максимальное значение src_created_at_ts из stg.bookings..." - ) - cur.execute("SELECT max(src_created_at_ts) FROM stg.bookings") - row = cur.fetchone() - last_ts = row[0] - - if last_ts is None: - logging.info( - "В stg.bookings пока нет данных — будет выполнена полная загрузка (full)." - ) - return None - - logging.info( - "Последний загруженный src_created_at_ts в stg.bookings: %s", last_ts - ) - # Возвращаем строку, чтобы её было проще использовать в шаблонах и XCom - return last_ts.isoformat() - - -def _extract_and_load_increment_via_pxf( - last_loaded_ts: str | None, - load_date: str, - batch_id: str, -) -> None: - """ - Читает дельту из stg.bookings_ext и вставляет её в stg.bookings. - - Логика: - - если last_loaded_ts is None → первая загрузка (full), - берём все данные из источника; - - иначе берём только записи, где src_created_at_ts > last_loaded_ts - и не позже конца учебного дня. - """ - logging.info( - "Загрузка инкремента через PXF: last_loaded_ts=%s, load_date=%s, batch_id=%s", - last_loaded_ts, - load_date, - batch_id, - ) - - with get_gp_conn() as conn, conn.cursor() as cur: - if last_loaded_ts is None: - # Полная загрузка: переносим все строки из внешней таблицы. - logging.info("Режим загрузки: full (первичная загрузка данных).") - cur.execute( - """ - INSERT INTO stg.bookings ( - book_ref, - book_date, - total_amount, - src_created_at_ts, - load_dttm, - batch_id - ) - SELECT - book_ref::text, - book_date::text, - total_amount::text, - book_date::timestamp, - now(), - %s - FROM stg.bookings_ext - """, - (batch_id,), - ) - else: - # Инкрементальная загрузка: берём только «новые» строки по окну времени. - logging.info("Режим загрузки: delta (инкрементальная загрузка).") - cur.execute( - """ - INSERT INTO stg.bookings ( - book_ref, - book_date, - total_amount, - src_created_at_ts, - load_dttm, - batch_id - ) - SELECT - book_ref::text, - book_date::text, - total_amount::text, - book_date::timestamp, - now(), - %s - FROM stg.bookings_ext - WHERE book_date > %s::timestamp - AND book_date <= (%s::date + INTERVAL '1 day') - """, - (batch_id, last_loaded_ts, load_date), - ) - - inserted = cur.rowcount if cur.rowcount not in (None, -1) else None - conn.commit() - - logging.info("Вставлено строк в stg.bookings: %s", inserted) - - -def _check_row_counts( - load_date: str, - last_loaded_ts: str | None, - batch_id: str, -) -> None: - """ - Проверяет, что количество строк из источника и в stg.bookings совпадает. - - Для наглядности считаем: - - количество строк в stg.bookings_ext за текущее окно; - - количество строк в stg.bookings с текущим batch_id. - """ - logging.info( - "Проверка количества строк за %s (last_loaded_ts=%s, batch_id=%s)", - load_date, - last_loaded_ts, - batch_id, - ) - - # Здесь мы специально не вытаскиваем rowcount из предыдущей задачи, - # а пересчитываем окно, чтобы показать связь DQ‑логики с бизнес-правилами. - with get_gp_conn() as conn, conn.cursor() as cur: - if last_loaded_ts is None: - # full: считаем все строки во внешней таблице - cur.execute("SELECT COUNT(*) FROM stg.bookings_ext") - src_count = cur.fetchone()[0] - else: - # delta: считаем строки только за текущий интервал - cur.execute( - """ - SELECT COUNT(*) - FROM stg.bookings_ext - WHERE book_date > %s::timestamp - AND book_date <= (%s::date + INTERVAL '1 day') - """, - (last_loaded_ts, load_date), - ) - src_count = cur.fetchone()[0] - - # Считаем количество строк, реально вставленных в stg.bookings в этом запуске - cur.execute( - """ - SELECT COUNT(*) - FROM stg.bookings - WHERE batch_id = %s - """, - (batch_id,), - ) - stg_count = cur.fetchone()[0] - - if src_count != stg_count: - raise ValueError( - "❌ Несовпадение количества строк при загрузке bookings: " - f"источник={src_count}, stg={stg_count}. " - "Проверьте логи задач extract_and_load_increment_via_pxf " - "и корректность окна инкремента." - ) - - logging.info( - "✅ Проверка количества строк пройдена: источник=%s, stg=%s", - src_count, - stg_count, - ) +BOOKINGS_CONN_ID = "bookings_db" +GREENPLUM_CONN_ID = "greenplum_conn" def _finish_summary() -> None: - """Логирует краткий итог выполнения DAG за один запуск.""" - logging.info("DAG bookings_to_gp_stage завершён.") + """ + Логирует краткий итог выполнения DAG за один запуск. + + Здесь можно добавить дополнительную агрегацию/логирование, + но для учебного примера достаточно простого сообщения. + """ + from logging import getLogger + + log = getLogger(__name__) + log.info("DAG bookings_to_gp_stage завершён. Подробности смотрите в логах задач.") with DAG( @@ -372,41 +50,40 @@ with DAG( tags=["demo", "bookings", "greenplum", "stg"], description="Учебный DAG: загрузка из bookings-db в stg.bookings (Greenplum)", ) as dag: - generate_bookings_day = PythonOperator( + # 1. Генерируем один учебный день в демо-БД bookings (идемпотентно) + generate_bookings_day = PostgresOperator( task_id="generate_bookings_day", - python_callable=_generate_bookings_day, - op_kwargs={"load_date": "{{ ds }}"}, + postgres_conn_id=BOOKINGS_CONN_ID, + sql="/sql/src/bookings_generate_day_if_missing.sql", + params={"load_date": "{{ ds }}"}, ) - get_last_loaded_ts = PythonOperator( - task_id="get_last_loaded_ts_from_gp", - python_callable=_get_last_loaded_ts_from_gp, - ) - - extract_and_load_increment = PythonOperator( - task_id="extract_and_load_increment_via_pxf", - python_callable=_extract_and_load_increment_via_pxf, - op_kwargs={ - "last_loaded_ts": "{{ ti.xcom_pull(task_ids='get_last_loaded_ts_from_gp') }}", + # 2. Загружаем инкремент из stg.bookings_ext в stg.bookings + load_bookings_to_stg = PostgresOperator( + task_id="load_bookings_to_stg", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="/sql/stg/bookings_load.sql", + params={ "load_date": "{{ ds }}", "batch_id": "{{ ds_nodash }}", }, ) - check_row_counts = PythonOperator( + # 3. Проверяем количество строк между источником и stg.bookings + check_row_counts = PostgresOperator( task_id="check_row_counts", - python_callable=_check_row_counts, - op_kwargs={ + postgres_conn_id=GREENPLUM_CONN_ID, + sql="/sql/stg/bookings_dq.sql", + params={ "load_date": "{{ ds }}", - "last_loaded_ts": "{{ ti.xcom_pull(task_ids='get_last_loaded_ts_from_gp') }}", "batch_id": "{{ ds_nodash }}", }, ) + # 4. Финальный лог/сводка finish_summary = PythonOperator( task_id="finish_summary", python_callable=_finish_summary, ) - generate_bookings_day >> get_last_loaded_ts >> extract_and_load_increment >> check_row_counts >> finish_summary - + generate_bookings_day >> load_bookings_to_stg >> check_row_counts >> finish_summary diff --git a/docker-compose.yml b/docker-compose.yml index e2b198d..d71335f 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -84,6 +84,7 @@ services: volumes: - ./airflow/dags:/opt/airflow/dags - ./airflow/requirements.txt:/opt/airflow/requirements.txt + - ./sql:/sql:ro - airflow_data:/opt/airflow/data depends_on: pgmeta: @@ -106,6 +107,7 @@ services: volumes: - ./airflow/dags:/opt/airflow/dags - ./airflow/requirements.txt:/opt/airflow/requirements.txt + - ./sql:/sql:ro - airflow_data:/opt/airflow/data depends_on: pgmeta: @@ -126,6 +128,7 @@ services: volumes: - ./airflow/dags:/opt/airflow/dags - ./airflow/requirements.txt:/opt/airflow/requirements.txt + - ./sql:/sql:ro - airflow_data:/opt/airflow/data command: > bash -lc " diff --git a/docs/internal/bookings_stg_design.md b/docs/internal/bookings_stg_design.md index 18474a2..75725b0 100644 --- a/docs/internal/bookings_stg_design.md +++ b/docs/internal/bookings_stg_design.md @@ -86,41 +86,45 @@ DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Gre ### 4.2. DAG для ежедневной загрузки -- `dag_id`: `bookings_to_gp_stage` (рабочее имя). +- `dag_id`: `bookings_to_gp_stage`. - Основные параметры: - - `load_date` (по умолчанию `{{ ds }}`) — учебный день, за который генерим и грузим данные. + - `load_date` (по умолчанию `{{ ds }}`) — учебный день, за который генерим и грузим данные; - подключения: - - `bookings_db_conn_id` — Airflow connection к `bookings-db`; - - `greenplum_conn_id` — Airflow connection к Greenplum. + - `bookings_db_conn_id` — Airflow connection к `bookings-db` (в коде DAG — `BOOKINGS_CONN_ID = "bookings_db"`); + - `greenplum_conn_id` — Airflow connection к Greenplum (`GREENPLUM_CONN_ID = "greenplum_conn"`). -Предлагаемая последовательность задач: +Последовательность задач (упрощённая, но отражающая те же шаги): 1. `generate_bookings_day` - - при необходимости создаёт таблицу `bookings.bookings` в `bookings-db` (одноразово); - - **идемпотентность**: проверяет, есть ли данные за `load_date` в `bookings.bookings`; - - если нет — генерирует данные за один учебный день (можно опираться на `make bookings-generate-day` и логику из `docs/internal/bookings_tz.md`); - - если да — ничего не делает, логирует «данные за этот день уже есть» и считается успешно выполненной. -2. `get_last_loaded_ts_from_gp` - - читает `max(src_created_at_ts)` из `stg.bookings`; - - если данных нет — возвращает `None` и режим загрузки `full`; - - результат (последняя метка и режим) передаёт через XCom. -3. `extract_and_load_increment_via_pxf` - - делает `INSERT INTO stg.bookings (...) SELECT ... FROM stg.bookings_ext WHERE src_created_at_ts > :last_ts AND src_created_at_ts <= :window_end`; - - в режиме `full` подставляет «минимальную» дату (или просто не фильтрует); - - заполняет `src_created_at_ts`, `load_dttm`, `batch_id`. -4. `check_row_counts` - - сравнивает количество строк, прочитанных из `stg.bookings_ext`, и количеством реально вставленных в `stg.bookings`; - - при расхождении падает с понятным сообщением «что посмотреть/проверить». -5. (опционально) `finish_summary` - - логирует итог: режим загрузки (full/delta), интервал `src_created_at_ts`, количество строк. + - PostgresOperator к `bookings-db`; + - выполняет скрипт `/sql/bookings/generate_day_if_missing.sql`; + - скрипт сам идемпотентно проверяет, есть ли данные за `load_date` в `bookings.bookings`: + - если день уже есть — выдаёт NOTICE и ничего не делает; + - если нет — вызывает генератор демобазы (логика как в `bookings/generate_next_day.sql`). +2. `load_bookings_to_stg` + - PostgresOperator к Greenplum; + - выполняет скрипт `/sql/stg/bookings_load.sql`; + - внутри SQL считается `max(src_created_at_ts)` по «старым» батчам и по нему строится окно инкремента: + - первая загрузка (full) — берём все строки из `stg.bookings_ext`; + - последующие загрузки — берём только записи, где `book_date` больше предыдущего максимума и не позже конца учебного дня; + - при вставке заполняются тех.колонки `src_created_at_ts`, `load_dttm`, `batch_id`. +3. `check_row_counts` + - PostgresOperator к Greenplum; + - выполняет скрипт `/sql/stg/bookings_dq.sql`; + - скрипт заново считает окно инкремента по тем же правилам, что и загрузка, и сравнивает: + - количество строк в `stg.bookings_ext` за окно, + - количество строк в `stg.bookings` для текущего `batch_id`; + - при расхождении выполняет `RAISE EXCEPTION` с понятным текстом ошибки. +4. `finish_summary` + - PythonOperator, который логирует итог выполнения DAG и напоминает, где смотреть детальные логи. -Эта последовательность служит учебным примером для менти: от генерации «операционных» данных до инкремента в сырой слой DWH. +Таким образом, вся бизнес‑логика инкремента и проверок живёт в SQL‑скриптах, а DAG отвечает за оркестрацию и подключение к нужным БД. Для менти это хороший пример разделения ответственности между SQL и Python. ## 5. Связь с остальными документами - `docs/internal/bookings_tz.md` — как готовится и генерируется источник `bookings-db`. - `docs/internal/pxf_bookings.md` — детали настройки PXF и внешней таблицы для чтения из `bookings-db`. -- `sql/ddl_gp.sql` — итоговый DDL для схемы `stg` и таблиц `stg.bookings_ext` / `stg.bookings` (применяется через `make ddl-gp`). +- `sql/stg/bookings_ddl.sql` — DDL для схемы `stg` и таблиц `stg.bookings_ext` / `stg.bookings` (подключается из `sql/ddl_gp.sql` и применяется через `make ddl-gp`). Дальнейшая модель DWH (слои ODS/DDS/DM, факт/измерения, SCD) должна быть спроектирована студентом по статье о моделировании данных, используя `stg.bookings` как входной слой. diff --git a/docs/internal/bookings_stg_readme.md b/docs/internal/bookings_stg_readme.md index fa00c83..a7fe356 100644 --- a/docs/internal/bookings_stg_readme.md +++ b/docs/internal/bookings_stg_readme.md @@ -21,7 +21,7 @@ _Внутренний файл, чтобы не забыть договорён - Стенд поднят: `make up`. - Демо‑БД bookings инициализирована: `make bookings-init`. - В Greenplum применён DDL (созданы схема `stg` и таблицы `stg.bookings_ext` / `stg.bookings`): - - `make ddl-gp` (использует `sql/ddl_gp.sql`). + - `make ddl-gp` (использует `sql/ddl_gp.sql`, который подтягивает `sql/stg/bookings_ddl.sql`). - В Airflow есть коннекты: - `greenplum_conn` (по умолчанию уже используется в helpers/greenplum.py); - `bookings_db` (Postgres к сервису `bookings-db`, если не хочется полагаться на ENV). @@ -29,23 +29,19 @@ _Внутренний файл, чтобы не забыть договорён ## 3. Последовательность задач в DAG - `generate_bookings_day`: - - проверяет наличие таблицы `bookings.bookings`; - - смотрит, есть ли строки за `load_date` (`book_date::date = load_date`); - - если день уже сгенерирован — ничего не делает (идемпотентность), просто логирует это; - - если нет — запускает генератор через DO‑блок (логика как в `bookings/generate_next_day.sql`). -- `get_last_loaded_ts_from_gp`: - - проверяет, что есть таблица `stg.bookings`; - - берёт `max(src_created_at_ts)` как последнюю загруженную метку; - - если NULL — значит в STG ещё нет данных, дальше идём в режим `full`. -- `extract_and_load_increment_via_pxf`: - - читает данные из `stg.bookings_ext` и вставляет в `stg.bookings`; - - при `last_loaded_ts is None` делает полную загрузку (full) — все строки; - - при delta берёт только строки с `book_date > last_loaded_ts` и `book_date <= load_date + 1 day`; - - приводит бизнес‑поля к `TEXT`, заполняет `src_created_at_ts`, `load_dttm`, `batch_id={{ ds_nodash }}`. + - PostgresOperator к `bookings-db`; + - выполняет SQL `/sql/src/bookings_generate_day_if_missing.sql`; + - скрипт сам проверяет, есть ли строки за `load_date` (`book_date::date = load_date`); + - если день уже сгенерирован — ничего не делает (идемпотентность), только пишет NOTICE. +- `load_bookings_to_stg`: + - PostgresOperator к Greenplum; + - выполняет SQL `/sql/stg/bookings_load.sql`; + - считает «старый» максимум `src_created_at_ts` (по предыдущим батчам) и грузит только новые строки из `stg.bookings_ext`, заполняя `src_created_at_ts`, `load_dttm`, `batch_id={{ ds_nodash }}`. - `check_row_counts`: - - пересчитывает количество строк в источнике (`stg.bookings_ext`) за текущее окно (full/delta); - - сравнивает с количеством строк в `stg.bookings` для текущего `batch_id`; - - если не сходится — падает с понятным сообщением и подсказкой «куда смотреть». + - PostgresOperator к Greenplum; + - выполняет SQL `/sql/stg/bookings_dq.sql`; + - за то же окно инкремента считает количество строк в источнике и в `stg.bookings` (по текущему `batch_id`); + - при расхождении делает `RAISE EXCEPTION` с понятным текстом ошибки. - `finish_summary`: - логирует итог выполнения DAG за одно срабатывание. @@ -78,4 +74,3 @@ _Внутренний файл, чтобы не забыть договорён - В будущем всё это нужно будет собрать в одну понятную историю для студента: - отдельный раздел «Учебный пример: bookings → stg → dwh»; - скриншоты DAG, примеры запросов и типичные ошибки. - diff --git a/sql/ddl_gp.sql b/sql/ddl_gp.sql index 6b72f3f..026e601 100644 --- a/sql/ddl_gp.sql +++ b/sql/ddl_gp.sql @@ -22,27 +22,6 @@ CREATE EXTERNAL TABLE public.ext_bookings_bookings ( LOCATION ('pxf://bookings.bookings?PROFILE=JDBC&SERVER=bookings-db') FORMAT 'CUSTOM' (formatter='pxfwritable_import'); --- Схема stg для сырого слоя DWH. -CREATE SCHEMA IF NOT EXISTS stg; - --- Внешняя таблица в схеме stg для чтения данных из bookings.bookings через PXF. -DROP EXTERNAL TABLE IF EXISTS stg.bookings_ext; -CREATE EXTERNAL TABLE stg.bookings_ext ( - book_ref CHAR(6), - book_date TIMESTAMP, - total_amount NUMERIC(10,2) -) -LOCATION ('pxf://bookings.bookings?PROFILE=JDBC&SERVER=bookings-db') -FORMAT 'CUSTOM' (formatter='pxfwritable_import'); - --- Внутренняя таблица stg.bookings — сырой слой, все бизнес-колонки как TEXT. -CREATE TABLE IF NOT EXISTS stg.bookings ( - book_ref TEXT, - book_date TEXT, - total_amount TEXT, - src_created_at_ts TIMESTAMP, - load_dttm TIMESTAMP NOT NULL DEFAULT now(), - batch_id TEXT NOT NULL -) -WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) -DISTRIBUTED BY (book_ref); +-- DDL для слоя stg по таблице bookings вынесен в отдельный файл. +-- Здесь подключаем его через psql \i, чтобы сохранить единый входной скрипт. +\i stg/bookings_ddl.sql diff --git a/sql/src/bookings_generate_day_if_missing.sql b/sql/src/bookings_generate_day_if_missing.sql new file mode 100644 index 0000000..d4b7692 --- /dev/null +++ b/sql/src/bookings_generate_day_if_missing.sql @@ -0,0 +1,59 @@ +-- Генерация учебного дня в демо-БД bookings. +-- Если данные за указанный день уже существуют, генерация не выполняется. + +DO $$ +DECLARE + v_has_day boolean; + v_load_date date := {{ params.load_date }}::date; + v_max_book_date timestamptz; + v_start_date timestamptz; + v_end_date timestamptz; + v_jobs integer := COALESCE(current_setting('bookings.jobs', true), '1')::integer; + v_init_days integer := COALESCE(current_setting('bookings.init_days', true), '1')::integer; + v_start_cfg text := COALESCE(current_setting('bookings.start_date', true), '2017-01-01'); +BEGIN + -- Проверяем, что демобаза установлена + IF to_regclass('bookings.bookings') IS NULL THEN + RAISE EXCEPTION 'Таблица bookings.bookings не найдена. Сначала выполните make bookings-init.'; + END IF; + + -- Проверяем, есть ли уже данные за нужный день + SELECT EXISTS ( + SELECT 1 + FROM bookings.bookings + WHERE book_date::date = v_load_date + ) + INTO v_has_day; + + IF v_has_day THEN + RAISE NOTICE 'Данные за % уже есть в bookings.bookings — генерация не требуется.', v_load_date; + RETURN; + END IF; + + -- Ищем последнюю сгенерированную дату + SELECT max(book_date) INTO v_max_book_date FROM bookings.bookings; + + IF v_max_book_date IS NULL THEN + -- База пустая: берём стартовую дату из конфигурации (или дефолтную) + v_start_date := date_trunc('day', v_start_cfg::timestamptz); + ELSE + -- Продолжаем с дня, следующего за максимальной датой + v_start_date := date_trunc('day', v_max_book_date) + interval '1 day'; + END IF; + + -- Первая генерация вызывает generate(), последующие — continue() + IF v_max_book_date IS NULL THEN + v_end_date := v_start_date + (v_init_days || ' days')::interval; + CALL generate(v_start_date, v_end_date, v_jobs); + ELSE + v_end_date := v_start_date + interval '1 day'; + CALL continue(v_end_date, v_jobs); + END IF; + + -- Ждём завершения фоновых джобов генератора, чтобы данные успели записаться + WHILE busy() LOOP + PERFORM pg_sleep(1); + END LOOP; + PERFORM dblink_disconnect(unnest(dblink_get_connections())); +END $$; + diff --git a/sql/stg/bookings_ddl.sql b/sql/stg/bookings_ddl.sql new file mode 100644 index 0000000..5d4b8ba --- /dev/null +++ b/sql/stg/bookings_ddl.sql @@ -0,0 +1,29 @@ +-- DDL для слоя STG по таблице bookings. +-- Используется как из общего скрипта ddl_gp.sql (через \i), +-- так и может выполняться отдельно при изменении схемы. + +-- Схема stg для сырого слоя DWH. +CREATE SCHEMA IF NOT EXISTS stg; + +-- Внешняя таблица в схеме stg для чтения данных из bookings.bookings через PXF. +DROP EXTERNAL TABLE IF EXISTS stg.bookings_ext; +CREATE EXTERNAL TABLE stg.bookings_ext ( + book_ref CHAR(6), + book_date TIMESTAMP, + total_amount NUMERIC(10,2) +) +LOCATION ('pxf://bookings.bookings?PROFILE=JDBC&SERVER=bookings-db') +FORMAT 'CUSTOM' (formatter='pxfwritable_import'); + +-- Внутренняя таблица stg.bookings — сырой слой, все бизнес-колонки как TEXT. +CREATE TABLE IF NOT EXISTS stg.bookings ( + book_ref TEXT, + book_date TEXT, + total_amount TEXT, + src_created_at_ts TIMESTAMP, + load_dttm TIMESTAMP NOT NULL DEFAULT now(), + batch_id TEXT NOT NULL +) +WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +DISTRIBUTED BY (book_ref); + diff --git a/sql/stg/bookings_dq.sql b/sql/stg/bookings_dq.sql new file mode 100644 index 0000000..879faf3 --- /dev/null +++ b/sql/stg/bookings_dq.sql @@ -0,0 +1,50 @@ +-- Проверка количества строк между источником stg.bookings_ext и стейджем stg.bookings. +-- Считаем строки за то же окно инкремента, что и при загрузке, +-- используя batch_id и "старый" максимум src_created_at_ts. + +DO $$ +DECLARE + v_load_date date := {{ params.load_date }}::date; + v_batch_id text := {{ params.batch_id | tojson }}::text; + v_prev_ts timestamp; + v_src_count bigint; + v_stg_count bigint; +BEGIN + -- Опорная метка: максимум src_created_at_ts среди предыдущих батчей + SELECT max(src_created_at_ts) + INTO v_prev_ts + FROM stg.bookings + WHERE batch_id <> v_batch_id + OR batch_id IS NULL; + + IF v_prev_ts IS NULL THEN + -- Полная загрузка: считаем все строки во внешней таблице + SELECT COUNT(*) INTO v_src_count FROM stg.bookings_ext; + ELSE + -- Инкремент: считаем только строки за текущее окно + SELECT COUNT(*) + INTO v_src_count + FROM stg.bookings_ext + WHERE book_date > v_prev_ts + AND book_date <= (v_load_date + INTERVAL '1 day'); + END IF; + + -- Считаем строки, реально вставленные в stg.bookings в этом батче + SELECT COUNT(*) + INTO v_stg_count + FROM stg.bookings + WHERE batch_id = v_batch_id; + + IF v_src_count <> v_stg_count THEN + RAISE EXCEPTION + 'Несовпадение количества строк при загрузке bookings: источник=%, stg=%. Проверьте окно инкремента и логи задач загрузки.', + v_src_count, + v_stg_count; + END IF; + + RAISE NOTICE + 'Проверка количества строк пройдена: источник=%, stg=%', + v_src_count, + v_stg_count; +END $$; + diff --git a/sql/stg/bookings_load.sql b/sql/stg/bookings_load.sql new file mode 100644 index 0000000..ea2908f --- /dev/null +++ b/sql/stg/bookings_load.sql @@ -0,0 +1,32 @@ +-- Загрузка инкремента из stg.bookings_ext в stg.bookings. +-- Окно инкремента определяется по src_created_at_ts: +-- берем строки, где book_date больше максимального src_created_at_ts +-- среди "старых" батчей и не позже конца учебного дня. + +INSERT INTO stg.bookings ( + book_ref, + book_date, + total_amount, + src_created_at_ts, + load_dttm, + batch_id +) +SELECT + book_ref::text, + book_date::text, + total_amount::text, + book_date::timestamp, + now(), + {{ params.batch_id | tojson }}::text +FROM stg.bookings_ext +WHERE book_date > COALESCE( + ( + SELECT max(src_created_at_ts) + FROM stg.bookings + WHERE batch_id <> {{ params.batch_id | tojson }}::text + OR batch_id IS NULL + ), + TIMESTAMP '1900-01-01 00:00:00' +) + AND book_date <= ({{ params.load_date }}::date + INTERVAL '1 day'); +