Files
airflow-greenplum/docs/internal/bookings_stg_design.md
T

11 KiB
Raw Blame History

Дизайн STG для bookings в Greenplum (черновик)

Внутренний документ для учебного стенда. Перед итоговой сдачей можно объединить с основной документацией.

1. Цель и общий контур

  • Источник: Postgres в контейнере bookings-db, база demo, таблица bookings.bookings (см. docs/internal/bookings_tz.md).
  • Цель: показываем путь данных от операционной БД до сырого слоя DWH в Greenplum.
  • В этом документе описываем только часть src (bookings-db) → STG (Greenplum). Слои ODS/DDS/DM студент проектирует сам по статье про моделирование DWH.

Логика на уровне слоёв (по статье):

  • src: оперативная система (bookings-db, схема bookings).
  • stg: сырой слой в Greenplum, максимально близкий к источнику, без бизнес‑логики.
  • дальше по заданию менти можно строить ods/dds/dm поверх STG.

2. Схема и таблицы в Greenplum

2.1. Схема

  • Используем одну схему stg в Greenplum.
  • В этой схеме будут:
    • внешняя таблица PXF для чтения из bookings-db;
    • внутренняя таблица STG для долговременного хранения «сырых» данных.

2.2. Внешняя таблица (PXF)

  • Имя таблицы: stg.bookings_ext.
  • Назначение: «окно» в исходную таблицу bookings.bookings в bookings-db через PXF (JDBC).
  • Типы колонок:
    • можем использовать «родные» типы из bookings.bookings (включая даты/числа);
    • задача внешней таблицы — корректно читать данные из источника, не заниматься приведением типов.

DDL будет добавлен в sql/ddl_gp.sql в блоке DDL для Greenplum (примерно по шаблону из docs/internal/pxf_bookings.md), с LOCATION ('pxf://bookings.bookings?PROFILE=JDBC&SERVER=bookings-db').

2.3. Внутренняя таблица STG

  • Имя таблицы: stg.bookings.
  • Назначение: хранить сырые данные из источника для последующей обработки (ODS/DDS/витрины).
  • Принципы моделирования:
    • все бизнес‑колонки из bookings.bookings храним как TEXT (как в примерах STG из статьи);
    • не делаем UPDATE/DELETE, только INSERT новых записей;
    • бизнес‑колонки по названию совпадают с источником (чтобы проще было маппить).

Технологические колонки:

  • src_created_at_ts TIMESTAMP — дата/время из источника, приведённая к TIMESTAMP:
    • используется как опорная колонка для инкрементальной загрузки;
    • заполняется из исходной даты/времени (created_at или аналог).
  • load_dttm TIMESTAMP NOT NULL DEFAULT now() — когда запись была загружена в STG.
  • batch_id TEXT NOT NULL — идентификатор «пачки» (например, {{ ds_nodash }} или run_id Airflow).
  • при необходимости позже можно добавить src_system TEXT, если появятся другие источники.

Колонки‑бизнес‑ключи (booking_id и т.п.) храним как TEXT. В слое DDS позже можно будет ввести суррогатные ключи и нормализовать модель под витрины.

3. Инкрементальная загрузка

3.1. Опорное поле для инкремента

  • Опорная колонка: src_created_at_ts (внутреннее имя в STG).
  • Источник значения:
    • берём из соответствующей колонки в bookings.bookings (например, book_date/created_at — будет уточнено при реализации);
    • при чтении через stg.bookings_ext приводим к TIMESTAMP.

3.2. Правила определения full/delta

  • При первом запуске, если таблица stg.bookings пуста:
    • считаем режим full — загружаем все строки из stg.bookings_ext.
  • При последующих запусках:
    • читаем max(src_created_at_ts) из stg.bookings за все предыдущие загрузки;
    • считаем, что нужно загрузить только строки, где src_created_at_ts больше этой максимальной метки и не позже конца текущего учебного дня.

Таким образом, вся логика инкремента «замкнута» на один техно‑столбец src_created_at_ts, который студент потом сможет использовать и на следующих слоях (например, в CDC‑логике).

4. DAG’и Airflow (логика на уровне задач)

4.1. DAG для DDL

  • dag_id: bookings_stg_ddl (рабочее имя).
  • Назначение: один раз (или при изменении схемы) создать необходимые объекты в Greenplum:
    • схему stg (если её ещё нет);
    • внешнюю таблицу stg.bookings_ext (PXF → bookings-db);
    • внутреннюю таблицу stg.bookings с текстовыми колонками и тех.полями.
  • Этот DAG не загружает данные, только подготавливает структуру.
  • Вся DDL‑логика (CREATE/ALTER/DROP) сосредоточена здесь; рабочие DAG’и занимаются только DML (INSERT/SELECT).

4.2. DAG для ежедневной загрузки

  • dag_id: bookings_to_gp_stage.
  • Основные параметры:
    • load_date (по умолчанию {{ ds }}) — учебный день, за который генерим и грузим данные;
    • подключения:
      • 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
    • 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 и напоминает, где смотреть детальные логи.

Таким образом, вся бизнес‑логика инкремента и проверок живёт в SQL‑скриптах, а DAG отвечает за оркестрацию и подключение к нужным БД. Для менти это хороший пример разделения ответственности между SQL и Python.

5. Связь с остальными документами

  • docs/internal/bookings_tz.md — как готовится и генерируется источник bookings-db.
  • docs/internal/pxf_bookings.md — детали настройки PXF и внешней таблицы для чтения из bookings-db.
  • 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 как входной слой.

6. Требования к читаемости и комментариям

  • DAG’и bookings_stg_ddl и bookings_to_gp_stage — это учебный материал для менти.
  • В коде DAG’ов должны быть:
    • понятные docstring на русском у всех функций и DAG;
    • короткие комментарии рядом с нетривиальной логикой (особенно вокруг инкремента и идемпотентности);
    • говорящие task_id и названия функций, отражающие их роль в процессе.
  • Цель: чтобы по одному только коду DAG студент мог восстановить архитектуру процесса и сопоставить её с теорией из статьи про моделирование DWH.