Доводка стилистики документации
This commit is contained in:
@@ -89,8 +89,9 @@ SELECT * FROM stg.bookings ORDER BY src_created_at_ts DESC LIMIT 10;
|
||||
Основные (для потока bookings → DWH):
|
||||
|
||||
- `bookings_stg_ddl` — создаёт `stg.bookings_ext`/`stg.bookings` и `stg.tickets_ext`/`stg.tickets` в Greenplum;
|
||||
- `bookings_to_gp_stage` — генерирует учебный день в `bookings-db`, грузит инкремент в `stg.bookings` и `stg.tickets`
|
||||
(через PXF), затем выполняет DQ‑проверки.
|
||||
- `bookings_stg_ddl` — создаёт/обновляет весь STG слой для bookings (9 таблиц: bookings, tickets, airports, airplanes,
|
||||
routes, seats, flights, segments, boarding_passes; включая внешние `*_ext` через PXF);
|
||||
- `bookings_to_gp_stage` — генерирует учебный день в `bookings-db`, затем загружает данные в STG и выполняет DQ‑проверки.
|
||||
|
||||
Вспомогательные (побочный трек с CSV):
|
||||
|
||||
|
||||
+2
-2
@@ -38,10 +38,10 @@
|
||||
- Проверить, что все 5 задач Success и логи содержат `Проверка пройдена`.
|
||||
|
||||
- DAG `bookings_to_gp_stage` (полная проверка цепочки bookings → Greenplum STG):
|
||||
- предварительно выполнить один раз: `make bookings-init` (установка демобазы `demo` в контейнере `bookings-db`) и `make ddl-gp` (создаёт `stg.bookings_ext` и `stg.bookings` в Greenplum);
|
||||
- предварительно выполнить один раз: `make bookings-init` (установка демобазы `demo` в контейнере `bookings-db`) и `make ddl-gp` (создаёт STG слой в Greenplum, включая внешние `*_ext` через PXF);
|
||||
- важно: DAG `bookings_stg_ddl` **не** создаёт базу `demo` в `bookings-db`; если вы делали `docker compose down -v` / `make clean`, `make bookings-init` обязателен;
|
||||
- включить DAG `bookings_to_gp_stage` и запустить `Trigger DAG`;
|
||||
- убедиться, что все задачи (`generate_bookings_day`, `load_bookings_to_stg`, `check_row_counts`, `finish_summary`) завершились со статусом Success;
|
||||
- убедиться, что все задачи завершились со статусом Success (включая загрузки справочников/транзакций и DQ);
|
||||
- при желании проверить данные: в `bookings-db` появился новый день, а в Greenplum в `stg.bookings` — строки с актуальным `batch_id` (см. пример запросов в разделе 5).
|
||||
|
||||
- (опционально, для менторов/разработчиков) Smoke-тест DAG через Airflow CLI без UI:
|
||||
|
||||
@@ -1,8 +1,11 @@
|
||||
from __future__ import annotations
|
||||
|
||||
"""
|
||||
Учебный DAG: создаёт схему stg и таблицы bookings_ext/bookings/tickets в Greenplum.
|
||||
Запускается вручную перед DAG загрузки bookings_to_gp_stage или после изменения DDL.
|
||||
Учебный DAG: создаёт/обновляет слой stg в Greenplum для демо-источника bookings.
|
||||
|
||||
Запускается вручную перед DAG загрузки `bookings_to_gp_stage` или после изменения DDL.
|
||||
Создаёт внешние таблицы PXF (`*_ext`) и внутренние таблицы STG (9 таблиц: bookings, tickets,
|
||||
airports, airplanes, routes, seats, flights, segments, boarding_passes).
|
||||
"""
|
||||
|
||||
from datetime import timedelta
|
||||
@@ -23,8 +26,8 @@ with DAG(
|
||||
catchup=False,
|
||||
template_searchpath="/sql",
|
||||
default_args=default_args,
|
||||
tags=["demo", "greenplum", "ddl", "bookings", "tickets", "stg"],
|
||||
description="Создаёт/обновляет stg.bookings_ext/bookings/tickets для учебного DAG",
|
||||
tags=["demo", "greenplum", "ddl", "bookings", "stg"],
|
||||
description="Учебный DDL DAG: создаёт/обновляет stg.* (PXF external + internal STG) для bookings",
|
||||
) as dag:
|
||||
apply_stg_bookings_ddl = PostgresOperator(
|
||||
task_id="apply_stg_bookings_ddl",
|
||||
@@ -82,7 +85,8 @@ with DAG(
|
||||
sql="stg/boarding_passes_ddl.sql",
|
||||
)
|
||||
|
||||
# Сначала создаются справочники, затем транзакционные таблицы (последовательно)
|
||||
# DDL применяем последовательно, чтобы порядок был понятным для новичков,
|
||||
# а ошибки — воспроизводимыми (в логах сразу видно, на каком объекте упали).
|
||||
(
|
||||
apply_stg_bookings_ddl
|
||||
>> apply_stg_tickets_ddl
|
||||
|
||||
@@ -13,6 +13,8 @@ from __future__ import annotations
|
||||
- генератор в демо-БД bookings добавляет следующий учебный день после max(book_date);
|
||||
- загрузка в Greenplum берёт все строки, появившиеся после предыдущих батчей;
|
||||
- `run_id` используется как метка запуска (в `batch_id`, в логах и DQ).
|
||||
|
||||
Важно: для инкрементальных таблиц «пустое окно инкремента» допустимо (это не ошибка).
|
||||
"""
|
||||
|
||||
from datetime import timedelta
|
||||
@@ -56,7 +58,7 @@ with DAG(
|
||||
template_searchpath="/sql",
|
||||
default_args=default_args,
|
||||
tags=["demo", "bookings", "greenplum", "stg"],
|
||||
description="Учебный DAG: загрузка из bookings-db в stg.bookings и stg.tickets (Greenplum)",
|
||||
description="Учебный DAG: загрузка из bookings-db в слой stg (Greenplum) + DQ проверки",
|
||||
) as dag:
|
||||
# 1. Генерируем один (или несколько стартовых) учебный день в демо-БД bookings
|
||||
generate_bookings_day = PostgresOperator(
|
||||
@@ -180,7 +182,7 @@ with DAG(
|
||||
sql="stg/boarding_passes_dq.sql",
|
||||
)
|
||||
|
||||
# 6. Финальный лог/сводка
|
||||
# Финальный лог/сводка
|
||||
finish_summary = PythonOperator(
|
||||
task_id="finish_summary",
|
||||
python_callable=_finish_summary,
|
||||
@@ -190,13 +192,15 @@ with DAG(
|
||||
generate_bookings_day >> load_bookings_to_stg >> check_row_counts
|
||||
check_row_counts >> load_tickets_to_stg >> check_tickets_dq
|
||||
|
||||
# Затем загружаются справочники (последовательная загрузка)
|
||||
# Затем загружаются справочники.
|
||||
# Для простоты (и более понятных логов для новичков) делаем это последовательно.
|
||||
# Если позже понадобится ускорить DAG, эти шаги можно распараллелить, сохранив зависимости.
|
||||
check_tickets_dq >> load_airports_to_stg >> check_airports_dq
|
||||
check_airports_dq >> load_airplanes_to_stg >> check_airplanes_dq
|
||||
check_airplanes_dq >> load_routes_to_stg >> check_routes_dq
|
||||
check_routes_dq >> load_seats_to_stg >> check_seats_dq
|
||||
|
||||
# Затем загружаются транзакции (последовательная загрузка)
|
||||
# Затем загружаются транзакции (тоже последовательно, по тем же причинам).
|
||||
check_seats_dq >> load_flights_to_stg >> check_flights_dq
|
||||
check_flights_dq >> load_segments_to_stg >> check_segments_dq
|
||||
check_segments_dq >> load_boarding_passes_to_stg >> check_boarding_passes_dq
|
||||
|
||||
@@ -29,7 +29,9 @@ make up
|
||||
make bookings-init
|
||||
```
|
||||
|
||||
3) В Greenplum созданы STG‑объекты `stg.bookings_ext`/`stg.bookings` и `stg.tickets_ext`/`stg.tickets` (выберите один вариант):
|
||||
3) В Greenplum созданы STG‑объекты (внешние `*_ext` через PXF и внутренние таблицы слоя `stg`)
|
||||
для всех таблиц потока: `bookings`, `tickets`, `airports`, `airplanes`, `routes`, `seats`, `flights`,
|
||||
`segments`, `boarding_passes` (выберите один вариант):
|
||||
|
||||
- учебный вариант: запустить DAG `bookings_stg_ddl` в Airflow UI;
|
||||
- технический шорткат: `make ddl-gp`.
|
||||
@@ -130,7 +132,7 @@ LIMIT 10;
|
||||
## Типичные ошибки
|
||||
|
||||
- `database "demo" does not exist`: демо‑БД не установлена → выполните `make bookings-init`.
|
||||
- Ошибки про `stg.bookings_ext`/`stg.bookings`: не применён DDL → запустите `bookings_stg_ddl` или `make ddl-gp`.
|
||||
- Ошибки про `stg.*`/`stg.*_ext`: не применён DDL → запустите `bookings_stg_ddl` или `make ddl-gp`.
|
||||
- Ошибки PXF (`protocol "pxf" does not exist`, connection refused): перезапустите `greenplum` и повторите DDL.
|
||||
Для технических деталей см. `docs/internal/pxf_bookings.md`.
|
||||
|
||||
|
||||
@@ -1,5 +1,8 @@
|
||||
# План ETL для загрузки `bookings.tickets` в STG слой
|
||||
|
||||
> Архивный документ: это рабочий план, который использовался при разработке.
|
||||
> Актуальная реализация потока — DAG `bookings_to_gp_stage` и SQL в `sql/stg/`.
|
||||
|
||||
**Ветка:** `chore/bookings-etl`
|
||||
**Цель:** Добавить загрузку таблицы `tickets` в STG слой Greenplum по аналогии с `bookings`
|
||||
|
||||
|
||||
@@ -8,6 +8,10 @@ _Внутренний документ для учебного стенда. П
|
||||
- Цель: показываем путь данных от операционной БД до сырого слоя DWH в Greenplum.
|
||||
- В этом документе описываем только часть `src (bookings-db) → STG (Greenplum)`. Слои ODS/DDS/DM студент проектирует сам по статье про моделирование DWH.
|
||||
|
||||
Примечание: в текущей версии стенда слой `stg` содержит не только `bookings`, но и остальные таблицы потока
|
||||
(`tickets`, `airports`, `airplanes`, `routes`, `seats`, `flights`, `segments`, `boarding_passes`).
|
||||
Ниже логика разобрана на примере `bookings`, потому что на нём проще показать принципы инкремента и батчей.
|
||||
|
||||
Логика на уровне слоёв (по статье):
|
||||
|
||||
- `src`: оперативная система (`bookings-db`, схема `bookings`).
|
||||
@@ -20,8 +24,8 @@ _Внутренний документ для учебного стенда. П
|
||||
|
||||
- Используем одну схему `stg` в Greenplum.
|
||||
- В этой схеме будут:
|
||||
- внешняя таблица PXF для чтения из `bookings-db`;
|
||||
- внутренняя таблица STG для долговременного хранения «сырых» данных.
|
||||
- внешние таблицы PXF `*_ext` для чтения из `bookings-db`;
|
||||
- внутренние таблицы STG `stg.*` для долговременного хранения «сырых» данных.
|
||||
|
||||
### 2.2. Внешняя таблица (PXF)
|
||||
|
||||
@@ -31,6 +35,8 @@ _Внутренний документ для учебного стенда. П
|
||||
- можем использовать «родные» типы из `bookings.bookings` (включая даты/числа);
|
||||
- задача внешней таблицы — корректно читать данные из источника, не заниматься приведением типов.
|
||||
|
||||
В текущей реализации аналогично созданы внешние таблицы `*_ext` и для остальных сущностей (см. `sql/stg/*_ddl.sql`).
|
||||
|
||||
DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Greenplum (примерно по шаблону из `docs/internal/pxf_bookings.md`), с `LOCATION ('pxf://bookings.bookings?PROFILE=JDBC&SERVER=bookings-db')`.
|
||||
|
||||
### 2.3. Внутренняя таблица STG
|
||||
@@ -46,7 +52,7 @@ DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Gre
|
||||
|
||||
- `src_created_at_ts TIMESTAMP` — дата/время из источника, приведённая к TIMESTAMP:
|
||||
- используется как опорная колонка для инкрементальной загрузки;
|
||||
- заполняется из исходной даты/времени (`created_at` или аналог).
|
||||
- заполняется из опорной даты/времени, принятой для конкретной сущности (например, для `bookings` — из `book_date`).
|
||||
- `load_dttm TIMESTAMP NOT NULL DEFAULT now()` — когда запись была загружена в STG.
|
||||
- `batch_id TEXT NOT NULL` — идентификатор «пачки» (например, `{{ ds_nodash }}` или `run_id` Airflow).
|
||||
- при необходимости позже можно добавить `src_system TEXT`, если появятся другие источники.
|
||||
@@ -59,8 +65,8 @@ DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Gre
|
||||
|
||||
- Опорная колонка: `src_created_at_ts` (внутреннее имя в STG).
|
||||
- Источник значения:
|
||||
- берём из соответствующей колонки в `bookings.bookings` (например, `book_date`/`created_at` — будет уточнено при реализации);
|
||||
- при чтении через `stg.bookings_ext` приводим к `TIMESTAMP`.
|
||||
- для `bookings` используем `book_date` из `bookings.bookings` (в демо‑БД это поле естественно “шагает” по дням);
|
||||
- при чтении через `stg.bookings_ext` приводим к `TIMESTAMP` и сохраняем в `stg.bookings.src_created_at_ts`.
|
||||
|
||||
### 3.2. Правила определения full/delta
|
||||
|
||||
@@ -79,8 +85,8 @@ DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Gre
|
||||
- `dag_id`: `bookings_stg_ddl` (реализован в `airflow/dags/bookings_stg_ddl.py`).
|
||||
- Назначение: один раз (или при изменении схемы) создать необходимые объекты в Greenplum:
|
||||
- схему `stg` (если её ещё нет);
|
||||
- внешнюю таблицу `stg.bookings_ext` (PXF → `bookings-db`);
|
||||
- внутреннюю таблицу `stg.bookings` с текстовыми колонками и тех.полями.
|
||||
- внешние таблицы `*_ext` и внутренние таблицы слоя `stg` для всех сущностей потока
|
||||
(см. `sql/stg/*_ddl.sql`).
|
||||
- Этот DAG не загружает данные, только подготавливает структуру.
|
||||
- Вся DDL‑логика (CREATE/ALTER/DROP) сосредоточена здесь; рабочие DAG’и занимаются только DML (INSERT/SELECT).
|
||||
|
||||
@@ -115,7 +121,8 @@ DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Gre
|
||||
- количество строк в `stg.bookings_ext` с `book_date` позже «старого» максимума,
|
||||
- количество строк в `stg.bookings` для текущего `batch_id`;
|
||||
- при расхождении выполняет `RAISE EXCEPTION` с понятным текстом ошибки.
|
||||
4. `finish_summary`
|
||||
4. Далее — загрузка и DQ для остальных таблиц потока (tickets, справочники, транзакции).
|
||||
5. `finish_summary`
|
||||
- PythonOperator, который логирует итог выполнения DAG и напоминает, где смотреть детальные логи.
|
||||
|
||||
Таким образом, вся бизнес‑логика инкремента и проверок живёт в SQL‑скриптах, а DAG отвечает за оркестрацию и подключение к нужным БД. Для менти это хороший пример разделения ответственности между SQL и Python.
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
# Временное ТЗ по блоку bookings (для текущей разработки)
|
||||
|
||||
_Этот файл внутренний, удалить перед итоговой сдачей._
|
||||
_Внутренний файл для наставника: поясняет, как устроен источник `bookings-db` и генерация данных. Студентам обычно не нужен._
|
||||
|
||||
- Контейнер `bookings-db` — отдельный сервис Postgres из `docker-compose.yml`, база по умолчанию `demo` (из upstream demodb), без переименований.
|
||||
- Доступ снаружи не блокируем (порт `5434` по умолчанию), чтобы позже читать через PXF и подключаться из Greenplum.
|
||||
|
||||
@@ -32,7 +32,7 @@ WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
|
||||
-- Обоснование: airplane_code — это уникальный идентификатор самолёта.
|
||||
-- Использование airplane_code обеспечивает:
|
||||
-- 1. Равномерное распределение данных по сегментам (airplane_code имеет высокую кардинальность)
|
||||
-- 2. Co-location данных airplanes и routes при JOIN по airplane_code
|
||||
-- 3. Co-location данных airplanes и seats при JOIN по airplane_code
|
||||
-- 4. Оптимизацию запросов, которые фильтруют или группируют по airplane_code
|
||||
-- 2. Коллокацию данных airplanes и seats при JOIN по airplane_code
|
||||
-- 3. Оптимизацию запросов, которые фильтруют или группируют по airplane_code
|
||||
-- Примечание: JOIN с таблицей routes (распределённой по route_no) может требовать motion.
|
||||
DISTRIBUTED BY (airplane_code);
|
||||
|
||||
@@ -15,7 +15,7 @@ BEGIN
|
||||
|
||||
IF v_src_count = 0 THEN
|
||||
RAISE EXCEPTION
|
||||
'В источнике airplanes_ext нет строк.';
|
||||
'В источнике airplanes_ext нет строк. Проверьте: bookings-db запущен, PXF работает, STG DDL применён (bookings_stg_ddl или make ddl-gp).';
|
||||
END IF;
|
||||
|
||||
-- Считаем строки, реально вставленные в stg.airplanes в этом батче
|
||||
|
||||
@@ -36,6 +36,6 @@ WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
|
||||
-- Обоснование: airport_code — это уникальный идентификатор аэропорта.
|
||||
-- Использование airport_code обеспечивает:
|
||||
-- 1. Равномерное распределение данных по сегментам (airport_code имеет высокую кардинальность)
|
||||
-- 2. Co-location данных airports и routes при JOIN по departure_airport/arrival_airport
|
||||
-- 3. Оптимизацию запросов, которые фильтруют или группируют по airport_code
|
||||
-- 2. Оптимизацию запросов, которые фильтруют или группируют по airport_code
|
||||
-- Примечание: JOIN с таблицей routes (распределённой по route_no) может требовать motion.
|
||||
DISTRIBUTED BY (airport_code);
|
||||
|
||||
@@ -15,7 +15,7 @@ BEGIN
|
||||
|
||||
IF v_src_count = 0 THEN
|
||||
RAISE EXCEPTION
|
||||
'В источнике airports_ext нет строк.';
|
||||
'В источнике airports_ext нет строк. Проверьте: bookings-db запущен, PXF работает, STG DDL применён (bookings_stg_ddl или make ddl-gp).';
|
||||
END IF;
|
||||
|
||||
-- Считаем строки, реально вставленные в stg.airports в этом батче
|
||||
|
||||
@@ -32,7 +32,7 @@ WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
|
||||
-- Ключ распределения: ticket_no
|
||||
-- Обоснование: ticket_no — это основной бизнес-ключ для билетов.
|
||||
-- Использование ticket_no обеспечивает:
|
||||
-- 1. Co-location данных boarding_passes и segments при JOIN по ticket_no
|
||||
-- 1. Коллокацию данных boarding_passes и segments при JOIN по ticket_no
|
||||
-- 2. Равномерное распределение данных по сегментам (ticket_no имеет высокую кардинальность)
|
||||
-- Примечание: stg.tickets распределена по book_ref, поэтому JOIN boarding_passes ↔ tickets по ticket_no может требовать motion.
|
||||
DISTRIBUTED BY (ticket_no);
|
||||
|
||||
@@ -29,7 +29,6 @@ WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
|
||||
-- Обоснование: book_ref — это уникальный идентификатор бронирования.
|
||||
-- Использование book_ref обеспечивает:
|
||||
-- 1. Равномерное распределение данных по сегментам (book_ref имеет высокую кардинальность)
|
||||
-- 2. Co-location данных bookings и tickets при JOIN по book_ref
|
||||
-- 2. Коллокацию данных bookings и tickets при JOIN по book_ref
|
||||
-- 3. Оптимизацию запросов, которые фильтруют или группируют по book_ref
|
||||
DISTRIBUTED BY (book_ref);
|
||||
|
||||
|
||||
@@ -17,7 +17,7 @@ BEGIN
|
||||
|
||||
IF v_src_count = 0 THEN
|
||||
RAISE EXCEPTION
|
||||
'В источнике routes_ext нет строк.';
|
||||
'В источнике routes_ext нет строк. Проверьте: bookings-db запущен, PXF работает, STG DDL применён (bookings_stg_ddl или make ddl-gp).';
|
||||
END IF;
|
||||
|
||||
-- Считаем строки, реально вставленные в stg.routes в этом батче
|
||||
|
||||
@@ -26,9 +26,6 @@ CREATE TABLE IF NOT EXISTS stg.seats (
|
||||
)
|
||||
WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
|
||||
-- Ключ распределения: airplane_code
|
||||
-- Обоснование: airplane_code обеспечивает co-location с таблицей airplanes.
|
||||
-- Использование airplane_code обеспечивает:
|
||||
-- 1. Co-location данных seats и airplanes при JOIN по airplane_code
|
||||
-- 2. Группировка мест по самолётам (в одном самолёте обычно много мест)
|
||||
-- 3. Оптимизацию запросов, которые фильтруют или группируют по airplane_code
|
||||
-- Обоснование: airplane_code обеспечивает коллокацию seats ↔ airplanes при JOIN по airplane_code
|
||||
-- (в MPP это уменьшает вероятность перераспределения данных / motion).
|
||||
DISTRIBUTED BY (airplane_code);
|
||||
|
||||
@@ -16,7 +16,7 @@ BEGIN
|
||||
|
||||
IF v_src_count = 0 THEN
|
||||
RAISE EXCEPTION
|
||||
'В источнике seats_ext нет строк.';
|
||||
'В источнике seats_ext нет строк. Проверьте: bookings-db запущен, PXF работает, STG DDL применён (bookings_stg_ddl или make ddl-gp).';
|
||||
END IF;
|
||||
|
||||
-- Считаем строки, реально вставленные в stg.seats в этом батче
|
||||
|
||||
@@ -30,7 +30,7 @@ WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
|
||||
-- Ключ распределения: ticket_no
|
||||
-- Обоснование: ticket_no — это основной бизнес-ключ для билетов.
|
||||
-- Использование ticket_no обеспечивает:
|
||||
-- 1. Co-location данных segments и boarding_passes при JOIN по ticket_no
|
||||
-- 1. Коллокацию данных segments и boarding_passes при JOIN по ticket_no
|
||||
-- 2. Равномерное распределение данных по сегментам (ticket_no имеет высокую кардинальность)
|
||||
-- Примечание: stg.tickets распределена по book_ref, поэтому JOIN segments ↔ tickets по ticket_no может требовать motion.
|
||||
DISTRIBUTED BY (ticket_no);
|
||||
|
||||
@@ -32,8 +32,7 @@ WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
|
||||
-- Ключ распределения: book_ref
|
||||
-- Обоснование: book_ref — это основной бизнес-ключ для бронирований.
|
||||
-- Использование book_ref обеспечивает:
|
||||
-- 1. Co-location данных tickets и bookings при JOIN по book_ref
|
||||
-- 1. Коллокацию данных tickets и bookings при JOIN по book_ref
|
||||
-- 2. Равномерное распределение данных по сегментам (book_ref имеет высокую кардинальность)
|
||||
-- 3. Оптимизацию запросов, которые фильтруют или группируют по book_ref
|
||||
DISTRIBUTED BY (book_ref);
|
||||
|
||||
|
||||
Reference in New Issue
Block a user