Files
airflow-greenplum/plans/stg_layer_implementation_plan.md
T
ddadmin 4d793a9c11 refactor(sql): заменен тип сжатия zlib на zstd для AO-таблиц
- Зачем:
  - zstd (level 1) является современным стандартом для Greenplum 6.0+, обеспечивая более высокую скорость декомпрессии и лучшее сжатие.
- Что:
  - обновлены все DDL стейджинга (STG) и базовых таблиц.
  - обновлена архитектурная документация (ADR-3) и планы реализации.
  - исправлены примеры кода в Airflow DAG и описании ETL.
- Проверка:
  - успешное выполнение CREATE TABLE с новыми параметрами в Greenplum 6.27.1.
2026-03-01 20:26:56 +03:00

30 KiB
Raw Blame History

План реализации STG слоя целиком

Статус: План готов к реализации Дата: 2026-01-17 Автор: Architect Mode

Обзор задачи

Согласно docs/internal/db_schema.md, STG слой реализован частично (2 из 9 таблиц: bookings, tickets). Необходимо реализовать оставшиеся 7 таблиц.

Стратегия загрузки данных

Тип таблиц Стратегия Обоснование
Справочники (airports, airplanes, routes, seats) Full load Маленький объём (<10K строк), простота реализации
Транзакции (flights, segments) Инкремент Больший объём; выбираем максимально естественное опорное поле
Транзакции (boarding_passes) Full snapshot В источнике строки создаются и обновляются со временем, простого инкремента без усложнений нет

Важно: под Full load в STG подразумеваем «сняли слепок и дописали в append-only таблицу с batch_id», а не TRUNCATE + INSERT. Это даёт простую идемпотентность (по batch_id) и сохраняет историю загрузок.

Оценка размера справочников

Справочник Примерный размер Оценка
airports ~5K-6K аэропортов Маленький
airplanes ~10 моделей самолётов Крошечный
seats ~1700-2000 записей Маленький
routes Ожидается несколько тысяч Маленький/Средний

Вывод: Все справочники очень маленькие (до нескольких тысяч строк). Даже если routes будет 5000-10000 строк - это всё равно минимальный объём для Greenplum.

Список таблиц для реализации

Таблица источника Таблица STG Тип данных Стратегия загрузки Опорное поле для инкремента
bookings.airports_data stg.airports Справочник Full -
bookings.airplanes_data stg.airplanes Справочник Full -
bookings.routes stg.routes Справочник Full -
bookings.seats stg.seats Справочник Full -
bookings.flights stg.flights Транзакции Инкремент scheduled_departure
bookings.segments stg.segments Транзакции Инкремент book_date (через tickets)
bookings.boarding_passes stg.boarding_passes Транзакции Full (snapshot) -

Паттерн реализации (на основе bookings/tickets)

Для каждой таблицы создаются 3 файла:

  1. sql/stg/{table}_ddl.sql - DDL для внешней и внутренней таблиц
  2. sql/stg/{table}_load.sql - Загрузка (full или инкремент)
  3. sql/stg/{table}_dq.sql - Проверки качества данных

Общая структура DDL файла

-- DDL для слоя STG по таблице {table}.
-- Используется как из общего скрипта ddl_gp.sql (через \i),
-- так и может выполняться отдельно при изменении схемы.

-- Схема stg для сырого слоя DWH.
CREATE SCHEMA IF NOT EXISTS stg;

-- Внешняя таблица в схеме stg для чтения данных из bookings.{table} через PXF.
DROP EXTERNAL TABLE IF EXISTS stg.{table}_ext;
CREATE EXTERNAL TABLE stg.{table}_ext (
    -- поля из источника
)
LOCATION ('pxf://bookings.{table}?PROFILE=JDBC&SERVER=bookings-db')
FORMAT 'CUSTOM' (formatter='pxfwritable_import');

-- Внутренняя таблица stg.{table} — сырой слой, все бизнес-колонки как TEXT.
CREATE TABLE IF NOT EXISTS stg.{table} (
    -- бизнес-колонки как TEXT
    src_created_at_ts TIMESTAMP,
    load_dttm TIMESTAMP NOT NULL DEFAULT now(),
    batch_id TEXT
)
WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1)
DISTRIBUTED BY ({distribution_key});

Примечание по PXF/JDBC типам: для «сложных» типов Postgres (например, jsonb, point, массивы, tstzrange) чаще всего проще и надёжнее объявлять колонки во внешней таблице как TEXT, чтобы избежать несовместимостей драйвера/маппинга типов. Внутренний STG всё равно хранит бизнес-поля как TEXT.

Общая структура LOAD файла (Full load для справочников)

-- Загрузка всех строк из stg.{table}_ext в stg.{table}.
-- Используем batch_id для отслеживания загрузки.

INSERT INTO stg.{table} (
    -- бизнес-колонки
    src_created_at_ts,
    load_dttm,
    batch_id
)
SELECT
    ext.{field}::text,
    now()::timestamp,
    '{{ run_id }}'::text
FROM stg.{table}_ext AS ext
WHERE NOT EXISTS (
    -- Защита от дублей в рамках одного batch_id
    SELECT 1
    FROM stg.{table} AS t
    WHERE t.batch_id = '{{ run_id }}'::text
        AND t.{pk} = ext.{pk}::text
);

-- Обновляем статистику для оптимизатора Greenplum
ANALYZE stg.{table};

Общая структура LOAD файла (Инкремент для транзакций)

-- Загрузка инкремента из stg.{table}_ext в stg.{table}.
-- Окно инкремента определяется по src_created_at_ts:
-- берём строки, где {increment_field} больше максимального src_created_at_ts
-- среди "старых" батчей; верхняя граница по дате не используется.

-- CTE для определения максимальной даты загрузки предыдущего батча
WITH max_batch_ts AS (
    SELECT COALESCE(MAX(src_created_at_ts), TIMESTAMP '1900-01-01 00:00:00') AS max_ts
    FROM stg.{table}
    WHERE batch_id <> '{{ run_id }}'::text
        OR batch_id IS NULL
)
INSERT INTO stg.{table} (
    -- бизнес-колонки
    src_created_at_ts,
    load_dttm,
    batch_id
)
SELECT
    ext.{field}::text,
    ext.{increment_field}::timestamp,
    now(),
    '{{ run_id }}'::text
FROM stg.{table}_ext AS ext
CROSS JOIN max_batch_ts AS mb
WHERE ext.{increment_field} > mb.max_ts
AND NOT EXISTS (
    SELECT 1
    FROM stg.{table} AS t
    WHERE t.batch_id = '{{ run_id }}'::text
        AND t.{pk} = ext.{pk}::text
);

-- Обновляем статистику для оптимизатора Greenplum
ANALYZE stg.{table};

Общая структура DQ файла

Для инкрементальных таблиц сравниваем окно инкремента (по src_created_at_ts) между источником и STG. Для full snapshot таблиц (справочники и boarding_passes) обычно достаточно сравнить общее количество строк в источнике с количеством строк, загруженных в текущий batch_id, плюс проверить дубликаты/NULL/ссылочную целостность.

-- Проверки качества данных для {table}

DO $$
DECLARE
    v_batch_id TEXT := '{{ run_id }}'::text;
    v_prev_ts TIMESTAMP;
    v_src_count BIGINT;
    v_stg_count BIGINT;
    v_dup_count BIGINT;
    v_null_count BIGINT;
    -- другие переменные для специфических проверок
BEGIN
    -- Опорная метка: максимум src_created_at_ts среди предыдущих батчей
    SELECT max(src_created_at_ts)
    INTO v_prev_ts
    FROM stg.{table}
    WHERE batch_id <> v_batch_id
        OR batch_id IS NULL;

    -- Источник: считаем строки во внешней таблице, которые вошли в окно инкремента
    SELECT COUNT(*)
    INTO v_src_count
    FROM stg.{table}_ext
    WHERE {increment_field} > COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00');

    IF v_src_count = 0 THEN
        RAISE EXCEPTION
            'В источнике {table}_ext нет строк для окна инкремента.';
    END IF;

    -- Считаем строки, реально вставленные в stg.{table} в этом батче
    SELECT COUNT(*)
    INTO v_stg_count
    FROM stg.{table}
    WHERE batch_id = v_batch_id;

    IF v_src_count <> v_stg_count THEN
        RAISE EXCEPTION
            'DQ FAILED: несовпадение количества строк. Источник: %, STG: %',
            v_src_count,
            v_stg_count;
    END IF;

    -- Проверка на дубликаты первичного ключа
    SELECT COUNT(*) - COUNT(DISTINCT {pk})
    INTO v_dup_count
    FROM stg.{table} AS t
    WHERE t.batch_id = v_batch_id;

    IF v_dup_count <> 0 THEN
        RAISE EXCEPTION
            'DQ FAILED: найдены дубликаты {pk} (batch_id=%): %',
            v_batch_id,
            v_dup_count;
    END IF;

    -- Проверка обязательных полей
    SELECT COUNT(*)
    INTO v_null_count
    FROM stg.{table} AS t
    WHERE t.batch_id = v_batch_id
        AND (t.{required_field} IS NULL OR t.{required_field} = '');

    IF v_null_count <> 0 THEN
        RAISE EXCEPTION
            'DQ FAILED: найдены строки с NULL в обязательных полях (batch_id=%): %',
            v_batch_id,
            v_null_count;
    END IF;

    RAISE NOTICE
        'DQ PASSED: {table} ок (batch_id=%): source=% stg=%',
        v_batch_id,
        v_src_count,
        v_stg_count;
END $$;

Детали реализации по таблицам

1. airports (справочник, full load)

Внешняя таблица: stg.airports_ext

  • Поля: airport_code, airport_name (JSONB), city (JSONB), country (JSONB), coordinates, timezone
  • PXF: pxf://bookings.airports_data?PROFILE=JDBC&SERVER=bookings-db

Внутренняя таблица: stg.airports

  • Бизнес-колонки как TEXT:
    • airport_code TEXT
    • airport_name TEXT
    • city TEXT
    • country TEXT
    • coordinates TEXT
    • timezone TEXT
  • Тех.колонки: src_created_at_ts, load_dttm, batch_id
  • Распределение: DISTRIBUTED BY (airport_code)
  • Обоснование: airport_code — это уникальный идентификатор аэропорта

Загрузка: Full (все строки при каждом запуске)

DQ проверки:

  • Count между источником и STG
  • Дубликаты airport_code
  • NULL обязательных полей (airport_code, airport_name, city, timezone)

2. airplanes (справочник, full load)

Внешняя таблица: stg.airplanes_ext

  • Поля: airplane_code, model (JSONB), range, speed
  • PXF: pxf://bookings.airplanes_data?PROFILE=JDBC&SERVER=bookings-db

Внутренняя таблица: stg.airplanes

  • Бизнес-колонки как TEXT:
    • airplane_code TEXT
    • model TEXT
    • range TEXT
    • speed TEXT
  • Тех.колонки: src_created_at_ts, load_dttm, batch_id
  • Распределение: DISTRIBUTED BY (airplane_code)
  • Обоснование: airplane_code — это уникальный идентификатор самолёта

Загрузка: Full

DQ проверки:

  • Count между источником и STG
  • Дубликаты airplane_code
  • NULL обязательных полей (airplane_code, model)

3. routes (справочник, full load)

Внешняя таблица: stg.routes_ext

  • Поля: route_no, validity (tstzrange), departure_airport, arrival_airport, airplane_code, days_of_week (int[]), scheduled_time, duration
  • PXF: pxf://bookings.routes?PROFILE=JDBC&SERVER=bookings-db

Внутренняя таблица: stg.routes

  • Бизнес-колонки как TEXT:
    • route_no TEXT
    • validity TEXT
    • departure_airport TEXT
    • arrival_airport TEXT
    • airplane_code TEXT
    • days_of_week TEXT
    • scheduled_time TEXT
    • duration TEXT
  • Тех.колонки: src_created_at_ts, load_dttm, batch_id
  • Распределение: DISTRIBUTED BY (route_no)
  • Обоснование: route_no — логический идентификатор маршрута; он нужен для JOIN с flights по route_no

Загрузка: Full

DQ проверки:

  • Count между источником и STG
  • Дубликаты (route_no, validity)
  • NULL обязательных полей (route_no, departure_airport, arrival_airport, airplane_code)
  • Ссылочная целостность на airports (departure_airport, arrival_airport)
  • Ссылочная целостность на airplanes (airplane_code)

4. seats (справочник, full load)

Внешняя таблица: stg.seats_ext

  • Поля: airplane_code, seat_no, fare_conditions
  • PXF: pxf://bookings.seats?PROFILE=JDBC&SERVER=bookings-db

Внутренняя таблица: stg.seats

  • Бизнес-колонки как TEXT:
    • airplane_code TEXT
    • seat_no TEXT
    • fare_conditions TEXT
  • Тех.колонки: src_created_at_ts, load_dttm, batch_id
  • Распределение: DISTRIBUTED BY (airplane_code)
  • Обоснование: co-location с airplanes для оптимизации JOIN

Загрузка: Full

DQ проверки:

  • Count между источником и STG
  • Дубликаты (airplane_code, seat_no)
  • NULL обязательных полей (airplane_code, seat_no, fare_conditions)
  • Ссылочная целостность на airplanes (airplane_code)

5. flights (транзакции, инкремент)

Внешняя таблица: stg.flights_ext

  • Поля: flight_id, route_no, status, scheduled_departure, scheduled_arrival, actual_departure, actual_arrival
  • PXF: pxf://bookings.flights?PROFILE=JDBC&SERVER=bookings-db

Внутренняя таблица: stg.flights

  • Бизнес-колонки как TEXT:
    • flight_id TEXT
    • route_no TEXT
    • status TEXT
    • scheduled_departure TEXT
    • scheduled_arrival TEXT
    • actual_departure TEXT
    • actual_arrival TEXT
  • Тех.колонки: src_created_at_ts, load_dttm, batch_id
  • src_created_at_ts = scheduled_departure
  • Распределение: DISTRIBUTED BY (flight_id)
  • Обоснование: flight_id — это уникальный идентификатор рейса

Загрузка: Инкремент по scheduled_departure

DQ проверки:

  • Count между источником и STG
  • Дубликаты flight_id
  • NULL обязательных полей (flight_id, route_no, status, scheduled_departure)
  • Ссылочная целостность на routes (route_no)

Примечание: flights.status/actual_* в источнике могут меняться со временем. Для учебного STG можно принять допущение "insert-only" (снимаем слепок на момент загрузки), либо усложнить и перезагружать скользящее окно по датам вылета.

6. segments (транзакции, инкремент)

Внешняя таблица: stg.segments_ext

  • Поля: ticket_no, flight_id, fare_conditions, price
  • PXF: pxf://bookings.segments?PROFILE=JDBC&SERVER=bookings-db

Внутренняя таблица: stg.segments

  • Бизнес-колонки как TEXT:
    • ticket_no TEXT
    • flight_id TEXT
    • fare_conditions TEXT
    • price TEXT
  • Тех.колонки: src_created_at_ts, load_dttm, batch_id
  • src_created_at_ts = берётся из bookings.book_date через JOIN с tickets
  • Распределение: DISTRIBUTED BY (ticket_no)
  • Обоснование: co-location с tickets для оптимизации JOIN

Загрузка: Инкремент по book_date (как в tickets)

DQ проверки:

  • Count между источником и STG
  • Дубликаты (ticket_no, flight_id)
  • NULL обязательных полей (ticket_no, flight_id, fare_conditions, price)
  • Ссылочная целостность на tickets (ticket_no)
  • Ссылочная целостность на flights (flight_id)

7. boarding_passes (транзакции, full snapshot)

Внешняя таблица: stg.boarding_passes_ext

  • Поля: ticket_no, flight_id, seat_no, boarding_no, boarding_time
  • PXF: pxf://bookings.boarding_passes?PROFILE=JDBC&SERVER=bookings-db

Внутренняя таблица: stg.boarding_passes

  • Бизнес-колонки как TEXT:
    • ticket_no TEXT
    • flight_id TEXT
    • seat_no TEXT
    • boarding_no TEXT
    • boarding_time TEXT
  • Тех.колонки: src_created_at_ts, load_dttm, batch_id
  • src_created_at_ts = now() (в этой таблице нет удобного поля для инкремента, потому что строки могут создаваться и обновляться со временем)
  • Распределение: DISTRIBUTED BY (ticket_no)
  • Обоснование: co-location с tickets/segments для оптимизации JOIN

Загрузка: Full snapshot (все строки при каждом запуске)

DQ проверки:

  • Count между источником и STG
  • Дубликаты (ticket_no, flight_id)
  • NULL обязательных полей (ticket_no, flight_id)
  • Ссылочная целостность на tickets (ticket_no)
  • Ссылочная целостность на segments (ticket_no, flight_id)

Примечание: в источнике boarding_passes строки сначала создаются при CHECK-IN (без boarding_time), а потом обновляются при BOARDING. Поэтому инкремент "по времени" без усложнений будет пропускать часть событий и/или изменения. Для учебного стенда самый стабильный вариант — снимать полный слепок.

Обновление существующих DAG

airflow/dags/bookings_stg_ddl.py

Добавить задачи для создания DDL новых таблиц:

apply_stg_airports_ddl = PostgresOperator(
    task_id="apply_stg_airports_ddl",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/airports_ddl.sql",
)

apply_stg_airplanes_ddl = PostgresOperator(
    task_id="apply_stg_airplanes_ddl",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/airplanes_ddl.sql",
)

apply_stg_routes_ddl = PostgresOperator(
    task_id="apply_stg_routes_ddl",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/routes_ddl.sql",
)

apply_stg_seats_ddl = PostgresOperator(
    task_id="apply_stg_seats_ddl",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/seats_ddl.sql",
)

apply_stg_flights_ddl = PostgresOperator(
    task_id="apply_stg_flights_ddl",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/flights_ddl.sql",
)

apply_stg_segments_ddl = PostgresOperator(
    task_id="apply_stg_segments_ddl",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/segments_ddl.sql",
)

apply_stg_boarding_passes_ddl = PostgresOperator(
    task_id="apply_stg_boarding_passes_ddl",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/boarding_passes_ddl.sql",
)

Зависимости:

  • Сначала создаются справочники (airports, airplanes, routes, seats)
  • Затем транзакционные таблицы (flights, segments, boarding_passes)

airflow/dags/bookings_to_gp_stage.py

Добавить задачи для загрузки новых таблиц:

# Загрузка справочников (full load)
load_airports_to_stg = PostgresOperator(
    task_id="load_airports_to_stg",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/airports_load.sql",
)

check_airports_dq = PostgresOperator(
    task_id="check_airports_dq",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/airports_dq.sql",
)

load_airplanes_to_stg = PostgresOperator(
    task_id="load_airplanes_to_stg",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/airplanes_load.sql",
)

check_airplanes_dq = PostgresOperator(
    task_id="check_airplanes_dq",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/airplanes_dq.sql",
)

load_routes_to_stg = PostgresOperator(
    task_id="load_routes_to_stg",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/routes_load.sql",
)

check_routes_dq = PostgresOperator(
    task_id="check_routes_dq",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/routes_dq.sql",
)

load_seats_to_stg = PostgresOperator(
    task_id="load_seats_to_stg",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/seats_load.sql",
)

check_seats_dq = PostgresOperator(
    task_id="check_seats_dq",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/seats_dq.sql",
)

# Загрузка транзакций (инкремент)
load_flights_to_stg = PostgresOperator(
    task_id="load_flights_to_stg",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/flights_load.sql",
)

check_flights_dq = PostgresOperator(
    task_id="check_flights_dq",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/flights_dq.sql",
)

load_segments_to_stg = PostgresOperator(
    task_id="load_segments_to_stg",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/segments_load.sql",
)

check_segments_dq = PostgresOperator(
    task_id="check_segments_dq",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/segments_dq.sql",
)

load_boarding_passes_to_stg = PostgresOperator(
    task_id="load_boarding_passes_to_stg",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/boarding_passes_load.sql",
)

check_boarding_passes_dq = PostgresOperator(
    task_id="check_boarding_passes_dq",
    postgres_conn_id=GREENPLUM_CONN_ID,
    sql="stg/boarding_passes_dq.sql",
)

Зависимости:

  • Сначала загружаются и проверяются bookings и tickets (уже есть)
  • Затем загружаются справочники (airports, airplanes, routes, seats)
  • Затем загружаются транзакции (flights, segments, boarding_passes)
  • В конце финальный лог

Обновление sql/ddl_gp.sql

Добавить подключение новых DDL файлов:

-- DDL для слоя stg по таблицам bookings и tickets вынесены в отдельные файлы.
-- Здесь подключаем их через psql \i, чтобы сохранить единый входной скрипт.
\i stg/bookings_ddl.sql
\i stg/tickets_ddl.sql

-- DDL для новых таблиц STG слоя
\i stg/airports_ddl.sql
\i stg/airplanes_ddl.sql
\i stg/routes_ddl.sql
\i stg/seats_ddl.sql
\i stg/flights_ddl.sql
\i stg/segments_ddl.sql
\i stg/boarding_passes_ddl.sql

Добавление тестов

Обновить tests/test_dags_smoke.py для проверки структуры обновлённых DAG:

def test_bookings_stg_ddl_dag_structure():
    dag = _load_dag("airflow.dags.bookings_stg_ddl")

    expected_tasks = {
        "apply_stg_bookings_ddl",
        "apply_stg_tickets_ddl",
        "apply_stg_airports_ddl",
        "apply_stg_airplanes_ddl",
        "apply_stg_routes_ddl",
        "apply_stg_seats_ddl",
        "apply_stg_flights_ddl",
        "apply_stg_segments_ddl",
        "apply_stg_boarding_passes_ddl",
    }
    assert expected_tasks.issubset(dag.task_dict.keys())

    # Проверка линейных зависимостей
    # ... (проверка зависимостей между задачами)
def test_bookings_to_gp_stage_dag_structure():
    dag = _load_dag("airflow.dags.bookings_to_gp_stage")

    expected_tasks = {
        "generate_bookings_day",
        "load_bookings_to_stg",
        "check_row_counts",
        "load_tickets_to_stg",
        "check_tickets_dq",
        "load_airports_to_stg",
        "check_airports_dq",
        "load_airplanes_to_stg",
        "check_airplanes_dq",
        "load_routes_to_stg",
        "check_routes_dq",
        "load_seats_to_stg",
        "check_seats_dq",
        "load_flights_to_stg",
        "check_flights_dq",
        "load_segments_to_stg",
        "check_segments_dq",
        "load_boarding_passes_to_stg",
        "check_boarding_passes_dq",
        "finish_summary",
    }
    assert expected_tasks.issubset(dag.task_dict.keys())

    # Проверка линейных зависимостей
    # ... (проверка зависимостей между задачами)

Обновление документации

Обновить статус в docs/internal/db_schema.md с "2 из 9" на "9 из 9".

Добавить описание новых таблиц в документацию.

Диаграмма потока данных STG слоя

graph TB
    subgraph Source[Source: bookings-db]
        B1[airports_data]
        B2[airplanes_data]
        B3[routes]
        B4[seats]
        B5[flights]
        B6[segments]
        B7[boarding_passes]
    end
    
    subgraph STG[STG Layer: Greenplum]
        S1[stg.airports]
        S2[stg.airplanes]
        S3[stg.routes]
        S4[stg.seats]
        S5[stg.flights]
        S6[stg.segments]
        S7[stg.boarding_passes]
    end
    
    B1 --> S1
    B2 --> S2
    B3 --> S3
    B4 --> S4
    B5 --> S5
    B6 --> S6
    B7 --> S7

Чек-лист реализации

  • Создать файлы DDL для новых таблиц (7 файлов)
    • sql/stg/airports_ddl.sql
    • sql/stg/airplanes_ddl.sql
    • sql/stg/routes_ddl.sql
    • sql/stg/seats_ddl.sql
    • sql/stg/flights_ddl.sql
    • sql/stg/segments_ddl.sql
    • sql/stg/boarding_passes_ddl.sql
  • Создать файлы LOAD для новых таблиц (7 файлов)
    • sql/stg/airports_load.sql
    • sql/stg/airplanes_load.sql
    • sql/stg/routes_load.sql
    • sql/stg/seats_load.sql
    • sql/stg/flights_load.sql
    • sql/stg/segments_load.sql
    • sql/stg/boarding_passes_load.sql
  • Создать файлы DQ для новых таблиц (7 файлов)
    • sql/stg/airports_dq.sql
    • sql/stg/airplanes_dq.sql
    • sql/stg/routes_dq.sql
    • sql/stg/seats_dq.sql
    • sql/stg/flights_dq.sql
    • sql/stg/segments_dq.sql
    • sql/stg/boarding_passes_dq.sql
  • Обновить DAG bookings_stg_ddl.py
  • Обновить DAG bookings_to_gp_stage.py
  • Обновить sql/ddl_gp.sql
  • Добавить тесты для новых DAG в tests/test_dags_smoke.py
  • Обновить документацию docs/internal/db_schema.md
  • Провести тестирование реализации

Примечания для реализации

  1. Именование файлов: Использовать {table}_ddl.sql, {table}_load.sql, {table}_dq.sql
  2. Ключи распределения: Выбирать ключи с высокой кардинальностью для равномерного распределения
  3. Co-location: Использовать одинаковые ключи распределения для связанных таблиц (tickets, segments, boarding_passes по ticket_no)
  4. Комментарии: Добавлять русскоязычные комментарии в SQL-файлы для студентов
  5. DQ проверки: Все проверки должны падать с RAISE EXCEPTION при ошибке
  6. Batch ID: Использовать {{ run_id }} для идентификации батча
  7. Защита от дублей: Использовать NOT EXISTS для предотвращения дублирования в рамках одного batch_id

Связанные документы