Ускорение выполнения аналитических SQL-запросов в Data Lakehouse на базе Trino + Iceberg + S3
Архитектура Data Lakehouse на базе Trino + Iceberg + S3 стала стандартом для современных аналитических платформ, где требуется гибкость, масштабируемость и работа с большими объемами данных без жёсткой привязки к проприетарным DWH.
Однако при переходе от тестовых наборов данных к многотерабайтным таблицам возникает проблема: даже относительно простые аналитические SQL-запросы начинают выполняться слишком долго или вовсе падают с ошибками из-за недостатка памяти на воркерах Trino.
В реальных проектах оптимизация производительности часто сводится не только к настройкам Trino, но и к комплексной переработке хранения и логики запросов. На практике применяются четыре основных направления:
- Оптимизация джоинов.
- Изменение структуры хранения данных.
- Партиционирование.
- Переписывание запросов под лимиты кластера.
Архитектура и узкие места
Компоненты
- Trino (PrestoSQL) — распределённый движок выполнения запросов.
- Iceberg — формат хранения с поддержкой транзакций, эволюции схемы и оптимизированных сканов.
- S3 — объектное хранилище, физический слой данных.
Основные узкие места
- Сетевые задержки при чтении большого количества мелких файлов с S3.
- Объём данных, попадающий в память воркеров — особенно при широких джоинах и агрегациях.
- Отсутствие эффективного фильтрации на уровне метаданных при плохом партиционировании.
- Неоптимальная организация Iceberg-таблиц (слишком мелкие или слишком крупные партиции).
- Неэффективный план выполнения из-за некорректных статистик в метаданных.
Оптимизация джоинов
Проблема
Trino по умолчанию использует hash join, который требует полной загрузки одной из таблиц в память. При джоинах «факт ↔ факт» или при соединении огромных таблиц это приводит к OOM (out of memory).
Решения
3.2.1. Broadcast join только для малых таблиц
SELECT /*+ BROADCAST(dim) */
f.date, f.sales, d.category
FROM fact_sales f
JOIN dim_products d ON f.product_id = d.id;- Когда использовать: размер таблицы dim < 10% от памяти воркера.
- Риск: если таблица окажется больше — Trino свалится.
Join reordering и фильтрация до джоина
Вместо:
SELECT * FROM big_fact f JOIN big_dim d ON f.id = d.id WHERE d.category = 'Electronics';
— лучше:
WITH filtered_dim AS (
SELECT * FROM big_dim WHERE category = 'Electronics'
)
SELECT *
FROM big_fact f
JOIN filtered_dim d ON f.id = d.id;- Плюс: фильтрация уменьшает объём данных до джоина.
- Риск: при неправильной статистике Trino может всё равно переставить порядок операций.
Bucket join
- Разбиение таблиц по одинаковым бакетам (bucketed tables) в Iceberg.
- Пример DDL:
CREATE TABLE big_fact (
id BIGINT,
...
)
WITH (
bucketed_by = ARRAY['id'],
bucket_count = 32
);- Плюс: Trino будет джоинить по бакетам, минимизируя shuffle.
- Риск: жёсткая привязка к числу бакетов, изменение требует полной переразметки данных.
Изменение структуры хранения
Укрупнение файлов
Мелкие Parquet-файлы (по 5–10 МБ) приводят к тысячам запросов к S3. Решение — компакция:
CALL system.optimize('my_table');- Оптимальный размер файла: 512 МБ – 1 ГБ для S3.
- Риск: слишком крупные файлы → длинные чтения и перерасход памяти.
Минимизация ненужных колонок
Iceberg хранит данные по колонкам, но Trino всё равно сканирует нужные файлы. Если таблица широка, стоит вынести редко используемые поля в отдельную таблицу.
Партиционирование
Правильный выбор ключа
- Дата — классический вариант (например, event_date).
- Сложный ключ — region + month или category + week для бизнес-аналитики.
Пример:
CREATE TABLE fact_sales (
date DATE,
region STRING,
sales DOUBLE
)
WITH (
partitioning = ARRAY['region', 'date_trunc(''month'', date)']
);
Динамическое партиционирование при записи
Trino + Iceberg позволяют задавать партиции прямо в INSERT ... SELECT.
Переписывание запросов
Избегание SELECT *
— всегда выбирайте только нужные поля, особенно при работе с широкими фактами.
Замена CTE на временные таблицы
Вместо:
WITH t1 AS (...) SELECT ... FROM t1 JOIN ...
— можно материализовать результат во временную Iceberg-таблицу и потом использовать.
Лимитирование оконных функций
При больших данных оконные функции (ROW_NUMBER, RANK) перегружают shuffle. Иногда их можно заменить агрегацией + join.
Пример комплексной оптимизации
До оптимизации:
SELECT d.region, SUM(f.sales) FROM big_fact f JOIN big_dim d ON f.id = d.id WHERE d.category = 'Electronics' AND f.date BETWEEN DATE '2024-01-01' AND DATE '2024-03-31' GROUP BY d.region;
- Выполнение: 17 минут, OOM на двух воркерах.
После оптимизации:
- Вынесена фильтрация по категории до джоина.
- Применено партиционирование по date.
- Укрупнены файлы.
- Переписан join на bucket join.
WITH filtered_dim AS (
SELECT id, region
FROM big_dim
WHERE category = 'Electronics'
)
SELECT d.region, SUM(f.sales)
FROM big_fact f
JOIN filtered_dim d ON f.id = d.id
WHERE f.date BETWEEN DATE '2024-01-01' AND DATE '2024-03-31'
GROUP BY d.region;- Выполнение: 2 мин 30 сек.
Как избежать рисков при оптимизации
- Тестируйте на продоподобных данных. Оптимизация на малых выборках может не сработать на полных объёмах.
- Следите за статистикой Iceberg. Обновляйте метаданные, иначе Trino примет неверный план.
- Используйте отдельный dev-кластер для экспериментов с bucket join и партиционированием.
- Делайте резервную копию таблицы перед массовой компакцией.
- Логируйте планы выполнения (EXPLAIN и EXPLAIN ANALYZE). Это поможет понять, где узкое место.
Оптимизация аналитических SQL-запросов в Lakehouse на Trino + Iceberg + S3 — это не только тюнинг движка, но и грамотная работа с данными:
- Джоины должны быть минимальны по объёму данных.
- Iceberg-таблицы должны иметь оптимальный размер файлов и партиционирование.
- Запросы должны быть переписаны так, чтобы нагрузка на shuffle и память была минимальной.
Применяя эти подходы в комплексе, можно сократить время выполнения тяжёлых запросов в 5–10 раз, а иногда и просто «спасти» их от падения по памяти.



