From 4d793a9c11a054af534835126b6d674686fea27a Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 1 Mar 2026 20:26:56 +0300 Subject: [PATCH] =?UTF-8?q?refactor(sql):=20=D0=B7=D0=B0=D0=BC=D0=B5=D0=BD?= =?UTF-8?q?=D0=B5=D0=BD=20=D1=82=D0=B8=D0=BF=20=D1=81=D0=B6=D0=B0=D1=82?= =?UTF-8?q?=D0=B8=D1=8F=20zlib=20=D0=BD=D0=B0=20zstd=20=D0=B4=D0=BB=D1=8F?= =?UTF-8?q?=20AO-=D1=82=D0=B0=D0=B1=D0=BB=D0=B8=D1=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - zstd (level 1) является современным стандартом для Greenplum 6.0+, обеспечивая более высокую скорость декомпрессии и лучшее сжатие. - Что: - обновлены все DDL стейджинга (STG) и базовых таблиц. - обновлена архитектурная документация (ADR-3) и планы реализации. - исправлены примеры кода в Airflow DAG и описании ETL. - Проверка: - успешное выполнение CREATE TABLE с новыми параметрами в Greenplum 6.27.1. --- airflow/dags/csv_to_greenplum.py | 2 +- docs/chore/bookings-etl.md | 4 ++-- docs/internal/architecture_review.md | 18 +++++++++--------- plans/stg_layer_implementation_plan.md | 2 +- sql/base/orders_ddl.sql | 2 +- sql/stg/airplanes_ddl.sql | 2 +- sql/stg/airports_ddl.sql | 2 +- sql/stg/boarding_passes_ddl.sql | 2 +- sql/stg/bookings_ddl.sql | 2 +- sql/stg/flights_ddl.sql | 2 +- sql/stg/routes_ddl.sql | 2 +- sql/stg/seats_ddl.sql | 2 +- sql/stg/segments_ddl.sql | 2 +- sql/stg/tickets_ddl.sql | 2 +- 14 files changed, 23 insertions(+), 23 deletions(-) diff --git a/airflow/dags/csv_to_greenplum.py b/airflow/dags/csv_to_greenplum.py index 8049b4f..f99eaa7 100644 --- a/airflow/dags/csv_to_greenplum.py +++ b/airflow/dags/csv_to_greenplum.py @@ -26,7 +26,7 @@ def _create_table() -> None: customer_id BIGINT NOT NULL, amount NUMERIC(12,2) NOT NULL ) - WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) + WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) DISTRIBUTED BY (order_id); """ with get_gp_conn() as conn, conn.cursor() as cur: diff --git a/docs/chore/bookings-etl.md b/docs/chore/bookings-etl.md index 755dd45..e910dcf 100644 --- a/docs/chore/bookings-etl.md +++ b/docs/chore/bookings-etl.md @@ -53,7 +53,7 @@ CREATE TABLE IF NOT EXISTS stg.tickets ( load_dttm TIMESTAMP NOT NULL DEFAULT now(), -- время загрузки в Greenplum batch_id TEXT -- идентификатор батча (из Airflow run_id) ) -WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) DISTRIBUTED BY (book_ref); -- распределение по ключу связи с bookings ``` @@ -123,7 +123,7 @@ CREATE TABLE IF NOT EXISTS stg.tickets ( load_dttm TIMESTAMP NOT NULL DEFAULT now(), batch_id TEXT ) -WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) -- Распределяем по book_ref, чтобы джойны tickets → bookings по book_ref были без motion. DISTRIBUTED BY (book_ref); diff --git a/docs/internal/architecture_review.md b/docs/internal/architecture_review.md index a290566..447e0f0 100644 --- a/docs/internal/architecture_review.md +++ b/docs/internal/architecture_review.md @@ -101,10 +101,10 @@ GP-специфичная best practice, которую забывают даж - 18 из 28 таблиц имеют неявный heap (нет `WITH`) — студент не видит, что выбор сделан - **Целевая раскладка по storage:** - **AO Column Store**: `dds.dim_calendar` (write-once, generate_series) - - **AO Row + zlib**: ODS snapshot-справочники (`airports`, `airplanes`, `routes`, `seats`) + - **AO Row + zstd**: ODS snapshot-справочники (`airports`, `airplanes`, `routes`, `seats`) — перевести загрузку с UPSERT на TRUNCATE+INSERT (честнее для full snapshot семантики) - - **AO Row + zlib**: `dds.dim_tariffs` (только INSERT, нет UPDATE) - - **AO Row + zlib**: `dm.route_performance` (full rebuild, по дизайну) + - **AO Row + zstd**: `dds.dim_tariffs` (только INSERT, нет UPDATE) + - **AO Row + zstd**: `dm.route_performance` (full rebuild, по дизайну) - **Heap (явный)**: ODS транзакционные (`bookings`, `tickets`, `flights`, `segments`, `boarding_passes`) — row-level UPDATE при SCD1 UPSERT - **Heap (явный)**: DDS измерения с UPDATE (`dim_airports`, `dim_airplanes`, @@ -228,7 +228,7 @@ GP-специфичная best practice, которую забывают даж ```sql -- 1. Собрать новую партицию во временную таблицу CREATE TABLE tmp_fact_20170102 (LIKE dds.fact_flight_sales) - WITH (appendonly=true, orientation=column, compresstype=zlib); + WITH (appendonly=true, orientation=column, compresstype=zstd); INSERT INTO tmp_fact_20170102 SELECT ... FROM ods... WHERE flight_date = '2017-01-02'; -- 2. Атомарно заменить партицию (без DELETE, без UPDATE) @@ -258,11 +258,11 @@ DROP TABLE tmp_fact_20170102; | Storage | Таблицы | Почему | |---------|---------|--------| -| **AO Column** zlib | `dds.dim_calendar` | Write-once (generate_series), никогда не обновляется. Колоночное хранение идеально для аналитических скан. | -| **AO Column** zlib | `dm.route_performance` | Full rebuild (TRUNCATE+INSERT), чисто аналитические чтения. | -| **AO Row** zlib | STG: все 9 таблиц | Уже реализовано. Append-only, иммутабельные батчи. | -| **AO Row** zlib | ODS snapshot: `airports`, `airplanes`, `routes`, `seats` | Полный snapshot каждый раз. Перевести загрузку с UPSERT на TRUNCATE+INSERT — честнее для семантики «текущий срез». | -| **AO Row** zlib | `dds.dim_tariffs` | Только INSERT новых тарифов, UPDATE не используется. | +| **AO Column** zstd | `dds.dim_calendar` | Write-once (generate_series), никогда не обновляется. Колоночное хранение идеально для аналитических скан. | +| **AO Column** zstd | `dm.route_performance` | Full rebuild (TRUNCATE+INSERT), чисто аналитические чтения. | +| **AO Row** zstd | STG: все 9 таблиц | Уже реализовано. Append-only, иммутабельные батчи. Примечание: используем **zstd (level 1)** вместо zlib, так как он обеспечивает более высокую скорость декомпрессии и лучшее сжатие в современных GP-кластерах (6.0+). | +| **AO Row** zstd | ODS snapshot: `airports`, `airplanes`, `routes`, `seats` | Полный snapshot каждый раз. Перевести загрузку с UPSERT на TRUNCATE+INSERT — честнее для семантики «текущий срез». | +| **AO Row** zstd | `dds.dim_tariffs` | Только INSERT новых тарифов, UPDATE не используется. | | **Heap** (явный) | ODS транзакционные: `bookings`, `tickets`, `flights`, `segments`, `boarding_passes` | Row-level UPDATE при SCD1 UPSERT. Heap обязателен. | | **Heap** (явный) | DDS измерения с UPDATE: `dim_airports`, `dim_airplanes`, `dim_passengers`, `dim_routes` | SCD1/SCD2 UPSERT с row-level UPDATE. | | **Heap** (явный) | `dds.fact_flight_sales` | UPDATE (is_boarded, seat_no меняются). | diff --git a/plans/stg_layer_implementation_plan.md b/plans/stg_layer_implementation_plan.md index 80a69a6..518d4ad 100644 --- a/plans/stg_layer_implementation_plan.md +++ b/plans/stg_layer_implementation_plan.md @@ -74,7 +74,7 @@ CREATE TABLE IF NOT EXISTS stg.{table} ( load_dttm TIMESTAMP NOT NULL DEFAULT now(), batch_id TEXT ) -WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) DISTRIBUTED BY ({distribution_key}); ``` diff --git a/sql/base/orders_ddl.sql b/sql/base/orders_ddl.sql index c1bfac6..f6b4caf 100644 --- a/sql/base/orders_ddl.sql +++ b/sql/base/orders_ddl.sql @@ -7,5 +7,5 @@ CREATE TABLE IF NOT EXISTS public.orders ( customer_id BIGINT NOT NULL, amount NUMERIC(12,2) NOT NULL ) -WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) DISTRIBUTED BY (order_id); diff --git a/sql/stg/airplanes_ddl.sql b/sql/stg/airplanes_ddl.sql index 5c85a09..f252983 100644 --- a/sql/stg/airplanes_ddl.sql +++ b/sql/stg/airplanes_ddl.sql @@ -27,7 +27,7 @@ CREATE TABLE IF NOT EXISTS stg.airplanes ( load_dttm TIMESTAMP NOT NULL DEFAULT now(), batch_id TEXT ) -WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) -- Ключ распределения: airplane_code -- Обоснование: airplane_code — это уникальный идентификатор самолёта. -- Использование airplane_code обеспечивает: diff --git a/sql/stg/airports_ddl.sql b/sql/stg/airports_ddl.sql index aa08084..f35a281 100644 --- a/sql/stg/airports_ddl.sql +++ b/sql/stg/airports_ddl.sql @@ -31,7 +31,7 @@ CREATE TABLE IF NOT EXISTS stg.airports ( load_dttm TIMESTAMP NOT NULL DEFAULT now(), batch_id TEXT ) -WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) -- Ключ распределения: airport_code -- Обоснование: airport_code — это уникальный идентификатор аэропорта. -- Использование airport_code обеспечивает: diff --git a/sql/stg/boarding_passes_ddl.sql b/sql/stg/boarding_passes_ddl.sql index eca5f55..5b432be 100644 --- a/sql/stg/boarding_passes_ddl.sql +++ b/sql/stg/boarding_passes_ddl.sql @@ -28,7 +28,7 @@ CREATE TABLE IF NOT EXISTS stg.boarding_passes ( load_dttm TIMESTAMP NOT NULL DEFAULT now(), batch_id TEXT ) -WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) -- Ключ распределения: ticket_no -- Обоснование: ticket_no — это основной бизнес-ключ для билетов. -- Использование ticket_no обеспечивает: diff --git a/sql/stg/bookings_ddl.sql b/sql/stg/bookings_ddl.sql index c9752ea..6c04f1e 100644 --- a/sql/stg/bookings_ddl.sql +++ b/sql/stg/bookings_ddl.sql @@ -24,7 +24,7 @@ CREATE TABLE IF NOT EXISTS stg.bookings ( load_dttm TIMESTAMP NOT NULL DEFAULT now(), batch_id TEXT NOT NULL ) -WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) -- Ключ распределения: book_ref -- Обоснование: book_ref — это уникальный идентификатор бронирования. -- Использование book_ref обеспечивает: diff --git a/sql/stg/flights_ddl.sql b/sql/stg/flights_ddl.sql index ece9a75..7a6d623 100644 --- a/sql/stg/flights_ddl.sql +++ b/sql/stg/flights_ddl.sql @@ -32,7 +32,7 @@ CREATE TABLE IF NOT EXISTS stg.flights ( load_dttm TIMESTAMP NOT NULL DEFAULT now(), batch_id TEXT ) -WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) -- Ключ распределения: flight_id -- Обоснование: flight_id — это уникальный идентификатор рейса. -- Использование flight_id обеспечивает: diff --git a/sql/stg/routes_ddl.sql b/sql/stg/routes_ddl.sql index 54f39d0..f10a78e 100644 --- a/sql/stg/routes_ddl.sql +++ b/sql/stg/routes_ddl.sql @@ -35,7 +35,7 @@ CREATE TABLE IF NOT EXISTS stg.routes ( load_dttm TIMESTAMP NOT NULL DEFAULT now(), batch_id TEXT ) -WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) -- Ключ распределения: route_no -- Обоснование: route_no — бизнес-идентификатор маршрута и часто используется в фильтрах/джойнах. -- Использование route_no обеспечивает: diff --git a/sql/stg/seats_ddl.sql b/sql/stg/seats_ddl.sql index 3dbd4fd..ae81bc4 100644 --- a/sql/stg/seats_ddl.sql +++ b/sql/stg/seats_ddl.sql @@ -24,7 +24,7 @@ CREATE TABLE IF NOT EXISTS stg.seats ( load_dttm TIMESTAMP NOT NULL DEFAULT now(), batch_id TEXT ) -WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) -- Ключ распределения: airplane_code -- Обоснование: airplane_code обеспечивает коллокацию seats ↔ airplanes при JOIN по airplane_code -- (в MPP это уменьшает вероятность перераспределения данных / motion). diff --git a/sql/stg/segments_ddl.sql b/sql/stg/segments_ddl.sql index e2d17ba..1d5273f 100644 --- a/sql/stg/segments_ddl.sql +++ b/sql/stg/segments_ddl.sql @@ -26,7 +26,7 @@ CREATE TABLE IF NOT EXISTS stg.segments ( load_dttm TIMESTAMP NOT NULL DEFAULT now(), batch_id TEXT ) -WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) -- Ключ распределения: ticket_no -- Обоснование: ticket_no — это основной бизнес-ключ для билетов. -- Использование ticket_no обеспечивает: diff --git a/sql/stg/tickets_ddl.sql b/sql/stg/tickets_ddl.sql index bc75f5f..9564ccb 100644 --- a/sql/stg/tickets_ddl.sql +++ b/sql/stg/tickets_ddl.sql @@ -28,7 +28,7 @@ CREATE TABLE IF NOT EXISTS stg.tickets ( load_dttm TIMESTAMP NOT NULL DEFAULT now(), batch_id TEXT ) -WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) -- Ключ распределения: book_ref -- Обоснование: book_ref — это основной бизнес-ключ для бронирований. -- Использование book_ref обеспечивает: