diff --git a/airflow/dags/bookings_to_gp_stage.py b/airflow/dags/bookings_to_gp_stage.py index 18e027b..88ad88c 100644 --- a/airflow/dags/bookings_to_gp_stage.py +++ b/airflow/dags/bookings_to_gp_stage.py @@ -10,6 +10,7 @@ docstring и комментарии помогают студенту понят """ import logging +import os from datetime import datetime, timedelta from airflow import DAG @@ -25,76 +26,341 @@ 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. + Готовит данные за указанный день в bookings-db. - Важно: функция должна быть идемпотентной: - если данные за load_date уже есть в исходной БД, - повторно пересобирать день не нужно. + Идея: + - если за load_date уже есть строки в bookings.bookings → ничего не делаем; + - если нет → запускаем генератор (аналог make bookings-generate-day). """ - logging.info( - "Генерация учебного дня в bookings-db за дату %s (заглушка)", load_date - ) - # TODO: реализовать проверку наличия дня в bookings.bookings - # и генерацию нового дня при его отсутствии. + 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). + Если данных ещё нет, возвращает None — это будет означать + режим полной загрузки (full). """ - logging.info( - "Чтение последнего загруженного src_created_at_ts из stg.bookings (заглушка)" - ) - # Пример будущей реализации: - # with get_gp_conn() as conn, conn.cursor() as cur: - # cur.execute("SELECT max(src_created_at_ts) FROM stg.bookings") - # row = cur.fetchone() - # return row[0] - return None + 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), - берём все данные за load_date и ранее; + берём все данные из источника; - иначе берём только записи, где src_created_at_ts > last_loaded_ts и не позже конца учебного дня. """ logging.info( - "Загрузка инкремента через PXF: last_loaded_ts=%s, load_date=%s (заглушка)", + "Загрузка инкремента через PXF: last_loaded_ts=%s, load_date=%s, batch_id=%s", last_loaded_ts, load_date, + batch_id, ) - # TODO: реализовать INSERT INTO stg.bookings (...) SELECT ... FROM stg.bookings_ext - # с учётом инкрементального окна по src_created_at_ts. + + 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) -> None: +def _check_row_counts( + load_date: str, + last_loaded_ts: str | None, + batch_id: str, +) -> None: """ Проверяет, что количество строк из источника и в stg.bookings совпадает. - Эта проверка должна помочь студенту увидеть пример простой DQ‑проверки - для инкрементальной загрузки. + Для наглядности считаем: + - количество строк в stg.bookings_ext за текущее окно; + - количество строк в stg.bookings с текущим batch_id. """ - logging.info("Проверка количества строк за %s (заглушка)", load_date) - # TODO: реализовать сравнение количества строк, - # например через SELECT COUNT(*) в источнике и в stg.bookings. + 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, + ) def _finish_summary() -> None: """Логирует краткий итог выполнения DAG за один запуск.""" - logging.info("DAG bookings_to_gp_stage завершён (пока только скелет).") + logging.info("DAG bookings_to_gp_stage завершён.") with DAG( @@ -123,13 +389,18 @@ with DAG( op_kwargs={ "last_loaded_ts": "{{ ti.xcom_pull(task_ids='get_last_loaded_ts_from_gp') }}", "load_date": "{{ ds }}", + "batch_id": "{{ ds_nodash }}", }, ) check_row_counts = PythonOperator( task_id="check_row_counts", python_callable=_check_row_counts, - op_kwargs={"load_date": "{{ ds }}"}, + op_kwargs={ + "load_date": "{{ ds }}", + "last_loaded_ts": "{{ ti.xcom_pull(task_ids='get_last_loaded_ts_from_gp') }}", + "batch_id": "{{ ds_nodash }}", + }, ) finish_summary = PythonOperator( @@ -138,3 +409,4 @@ with DAG( ) generate_bookings_day >> get_last_loaded_ts >> extract_and_load_increment >> check_row_counts >> finish_summary + diff --git a/docs/internal/bookings_stg_readme.md b/docs/internal/bookings_stg_readme.md new file mode 100644 index 0000000..fa00c83 --- /dev/null +++ b/docs/internal/bookings_stg_readme.md @@ -0,0 +1,81 @@ +# Мини‑README по учебному DAG bookings_to_gp_stage (черновик) + +_Внутренний файл, чтобы не забыть договорённости. Перед итоговой сдачей документацию по блоку bookings/STG нужно будет аккуратно собрать и переписать._ + +## 1. Что делает DAG + +- DAG `bookings_to_gp_stage` показывает учебный поток: + - источник: демо‑БД `bookings-db` (Postgres, схема `bookings`, таблица `bookings.bookings`); + - при каждом запуске генерируется один учебный день данных (идемпотентно); + - данные из `bookings.bookings` переливаются в сырой слой `stg.bookings` в Greenplum через PXF‑внешнюю таблицу `stg.bookings_ext`. +- Слой `stg` задуман как «сырой»: + - все бизнес‑колонки (`book_ref`, `book_date`, `total_amount`) хранятся как `TEXT`; + - есть тех.колонки `src_created_at_ts`, `load_dttm`, `batch_id`. + +Подробный дизайн описан в `docs/internal/bookings_stg_design.md`. + +## 2. Что нужно, чтобы DAG завёлся + +Минимальные предпосылки: + +- Стенд поднят: `make up`. +- Демо‑БД bookings инициализирована: `make bookings-init`. +- В Greenplum применён DDL (созданы схема `stg` и таблицы `stg.bookings_ext` / `stg.bookings`): + - `make ddl-gp` (использует `sql/ddl_gp.sql`). +- В Airflow есть коннекты: + - `greenplum_conn` (по умолчанию уже используется в helpers/greenplum.py); + - `bookings_db` (Postgres к сервису `bookings-db`, если не хочется полагаться на ENV). + +## 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 }}`. +- `check_row_counts`: + - пересчитывает количество строк в источнике (`stg.bookings_ext`) за текущее окно (full/delta); + - сравнивает с количеством строк в `stg.bookings` для текущего `batch_id`; + - если не сходится — падает с понятным сообщением и подсказкой «куда смотреть». +- `finish_summary`: + - логирует итог выполнения DAG за одно срабатывание. + +## 4. Как этим пользоваться студенту (черновой сценарий) + +1. Поднять стенд и подготовить источники: + - `make up` + - `make airflow-init` + - `make bookings-init` + - `make ddl-gp` +2. Открыть Airflow UI (`http://localhost:8080`) и включить DAG `bookings_to_gp_stage`. +3. Вызвать `Trigger` DAG с конкретной датой (например, `2017-01-01` → зависит от стартовой конфигурации демобазы). +4. Посмотреть: + - в `bookings-db` ― что появился день с бронированиями; + - в Greenplum (`make gp-psql`) — данные в `stg.bookings`: + - `SELECT * FROM stg.bookings LIMIT 10;` + - `SELECT src_created_at_ts, load_dttm, batch_id FROM stg.bookings ORDER BY src_created_at_ts DESC LIMIT 10;` +5. Перезапустить DAG для следующей даты и увидеть, что: + - генерация в `bookings.bookings` идёт по одному дню вперёд; + - в `stg.bookings` появляются только новые записи (delta). + +## 5. Примечания «на потом» + +- Текущая документация по блоку bookings/STG разбросана: + - `README.md` (общий обзор стенда), + - `docs/internal/bookings_tz.md` (источник bookings), + - `docs/internal/pxf_bookings.md` (PXF), + - `docs/internal/bookings_stg_design.md` (дизайн STG), + - этот файл (мини‑README по DAG). +- В будущем всё это нужно будет собрать в одну понятную историю для студента: + - отдельный раздел «Учебный пример: bookings → stg → dwh»; + - скриншоты DAG, примеры запросов и типичные ошибки. + diff --git a/sql/ddl_gp.sql b/sql/ddl_gp.sql index 2788c57..6b72f3f 100644 --- a/sql/ddl_gp.sql +++ b/sql/ddl_gp.sql @@ -21,3 +21,28 @@ 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);