Модуль 3. Развёртывание и конфигурация StarRocks
Методологическое понимание развёртывания
Развёртывание StarRocks — это не просто «скачать бинарник и запустить сервис».
В распределённой аналитической системе мы должны сразу решить четыре ключевых вопроса:
- Где и на чём размещать — bare-metal, виртуалки, контейнеры, Kubernetes.
- Как масштабировать — scale-up (увеличиваем ресурсы ноды) или scale-out (добавляем новые ноды).
- Как обеспечить отказоустойчивость — сколько FE и BE нужно, как их реплицировать.
- Как конфигурировать под нагрузку — ingestion vs BI-запросы, батчи vs real-time.
Методология здесь такая:
- Для теста — упрощённая установка, чтобы быстро пощупать возможности.
- Для продакшна — развёртывание с закладкой роста на 12–18 месяцев, с HA и резервами по CPU/RAM/Storage.
Аппаратные и системные требования
Минимальные рекомендации для продакшна
-
Frontend (FE):
- CPU: 4–8 vCPU (лучше 8)
- RAM: ≥16 GB
- Диск: SSD 100 GB
- Сеть: ≥1 Gbps (лучше 10 Gbps)
- Backend (BE):
- CPU: 8–32 vCPU
- RAM: ≥64 GB
- Диск: NVMe SSD 1–4 TB
- Сеть: ≥10 Gbps
- ОС: Linux (CentOS 7+, Ubuntu 20.04+)
- Java: OpenJDK 11+ (для FE)
Пример топологии на 1 млрд строк / день
- 3 FE (HA)
- 6 BE (по 128 GB RAM)
- Хранилище NVMe, RAID10
Варианты развёртывания
-
Single-node (Standalone) — для обучения, теста функций.
- FE и BE на одной машине.
- Простой запуск через start_fe.sh и start_be.sh.
- Ограничения: нет отказоустойчивости, падение = полный простой.
- Multi-node (Prod Cluster) — для боевой эксплуатации.
- Минимум 3 FE (1 лидер + 2 фолловера).
- Минимум 3 BE (с репликацией сегментов).
- HA за счёт распределения ролей.
- Helm Charts.
- StatefulSets для FE/BE.
- Динамическое добавление BE.
- Kubernetes / Cloud — для гибкого масштабирования.
Пошаговое развёртывание Multi-node кластера
Пример для bare-metal:
Шаг 1: Подготовка ОС
# Отключить swap swapoff -a # Настроить лимиты ulimit -n 65535 # Установить Java apt install openjdk-11-jdk
Шаг 2: Установка Frontend
tar zxvf StarRocks-x.y.z.tar.gz cd StarRocks/fe sh start_fe.sh --daemon
Лидер задаётся первым FE, остальные подключаются как followers:
ALTER SYSTEM ADD FOLLOWER "fe2_host:9010"; ALTER SYSTEM ADD FOLLOWER "fe3_host:9010";
Шаг 3: Установка Backend
cd StarRocks/be sh start_be.sh --daemon
Добавление BE через FE:
ALTER SYSTEM ADD BACKEND "be1_host:9050"; ALTER SYSTEM ADD BACKEND "be2_host:9050"; ALTER SYSTEM ADD BACKEND "be3_host:9050";
Шаг 4: Проверка кластера
SHOW PROC '/frontends'; SHOW PROC '/backends';
Конфигурационные параметры (ключевые)
FE:
- priority_networks — IP-пул для внутреннего трафика.
- edit_log_port — порт синхронизации FE.
- query_timeout — таймаут запросов (в сек).
- max_query_instances — лимит параллельных инстансов на BE.
BE:
- storage_root_path — путь к сегментам.
- be_port — основной порт RPC.
- brpc_max_body_size — максимальный размер передачи данных.
- disable_auto_compaction — авто-компакшн (для real-time нагрузки лучше настраивать вручную).
Практические кейсы конфигурации
Кейс 1. BI на 300 аналитиков
- Задача: обеспечить 300 параллельных сессий Tableau.
-
Решение:
- 4 FE (2 активных, 2 пассивных).
- 8 BE (по 96 GB RAM, hash distribution).
- max_query_instances=64 для балансировки нагрузки.
- Риск: BI-отчёты стреляют в сырые данные.
- Защита: только MVs в публичной схеме.
Кейс 2. Real-time ingestion из Kafka
- Задача: загружать 2 млн записей в минуту без лагов.
-
Решение:
- Routine Load с batch.size=50K.
- BE с NVMe RAID10.
- Auto-compaction ночью.
- Риск: накопление старых партиций.
- Защита: TTL на старые сегменты + pre-aggregation.
Риски при развёртывании и как их избежать
|
Риск |
Симптом |
Защита |
|---|---|---|
|
Один FE без резервов |
Падение = остановка кластера |
Минимум 3 FE в HA |
|
Неравномерная нагрузка на BE |
Одни BE перегружены |
Hash-дистрибуция по ключу, балансировка |
|
Заполнение диска BE |
Ошибки записи |
Мониторинг storage_root_path, TTL |
|
Падение ingestion при пике |
BI видит неполные данные |
Backpressure в Kafka/Flink |
|
Конфликты портов |
Ошибки подключения |
Явно указывать порты в конфиге |
Методологические советы
- Не экономить на FE — это мозг кластера, при его перегрузке падает всё.
- Отдельные сети для ingestion и клиентских запросов.
- Мониторинг с первого дня (Prometheus + Grafana).
- Тестировать ingestion под пиковую нагрузку до запуска.
- Держать документацию конфигов в Git, чтобы воспроизводить настройки.
Цели сайзинга (что закладываем)
- объем данных (исторический + прирост) и коэффициент репликации;
- профиль запросов (ад-hoc против типовых BI, доля тяжелых join/агрегаций, SLA по латентности);
- конкуренция (одновременных пользователей / активных запросов);
- ingest (батч/поток, пиковая скорость);
- запас по росту (обычно 12–18 месяцев);
- требования к отказоустойчивости (N FE, фактор репликации, многозонность).
Роли и базовые «нормы»
FE (Frontend) — планирование, метаданные, координация.
BE (Backend) — хранение сегментов и выполнение запросов.
- FE: 3 узла для HA (1 leader + 2 follower). Часто CPU «холодные», но чувствительны к памяти JVM и быстрым SSD для журнала.
- BE: основная мощность. Масштабируем по ядрам/памяти/дискам.
Правила больших пальцев (стартовые):
- Ядра: 16–32 vCPU на один BE для общего BI; 32–64 vCPU — для тяжелых join/RT.
- Память: 2–3 GB RAM на 1 vCPU BE (часто комфортно 256 GB RAM на 32 vCPU).
- Диск: NVMe, суммарный sequential ≥2–4 GB/s на BE; IOPS важны для компакшна и слияний.
- Сеть: внутри кластера 25 GbE (минимум 10 GbE), отдельно — ingress (ETL/Kafka) и client (BI).
- Репликация сегментов: RF=3 в проде (RF=2 допустим в dev/низкокритичных витринах).
- Запас по диску: не менее 30% свободного для компакшна и ребалансировки.
2) Хранилище: как считать объем
Колонки сжимаются. Для смешанных датасетов разумно брать компрессию 3–4× (консервативно 3.5×). Плюс метаданные/индексы/служебные файлы.
Формула (на 12–18 месяцев):
Storage_total = (Raw_hist + Raw_new * Horizon) / Compression * RF * (1 + Overhead) * (1 + Headroom)
Где:
- Compression ≈ 3.5
- RF = 3 (прод)
- Overhead ≈ 0.15 (индексы, MV, словари)
- Headroom = 0.3 (30% свободного места)
Если активно используете MV/предагрегации, добавляйте ещё 0.2–0.5× от сжатого объема «на жизнь» MVs (или считайте каждую крупную MV отдельно как еще одну таблицу).
TTL/ретеншн: на сырых/детальных партициях задайте срок хранения (например, 90 дней), всё старше — агрегируйте в более грубые витрины: это резко снижает «хвост» хранения.
3) Память и CPU: как оценить под запросы
Потребление памяти определяется:
- шириной набора (число столбцов),
- кардинальностью ключей (группировки/джоины),
- размером «рабочих наборов» при shuffle/join/aggregate,
- степенью параллелизма (сколько фрагментов исполняется одновременно).
Оценка памяти «на сложный запрос»:
- Простые BI-агрегации (без многократных больших join): 0.5–1.5 GB на BE-инстанс.
- Тяжелые запросы (несколько join’ов 100M+ строк, оконные функции): 2–6 GB на BE-инстанс.
- PK-таблицы с активным upsert и частыми компакшнами — держат больше временных структур.
Оценка одновременно активных запросов (Concurrency):
- Одновременные пользователи BI ≠ активные тяжелые запросы. Обычно активных тяжелых = 10–20% от одновременных BI-пользователей.
- Планируйте 1–2 активных тяжёлых запроса на 8–16 vCPU каждого BE, чтобы не ловить пики GC/pressure.
Правило для CPU:
- Для смешанного профиля: ~1 vCPU BE на 2–3 активных лёгких запроса или 1 vCPU на 0.5–1 тяжёлый (в пересчете на параллельные фрагменты).
4) Ingest: батчи и real-time
- Batch (Broker/Stream Load): считать пиковую суточную/часовую загрузку; выдерживать окна.
-
Routine Load / Kafka: считать пиковый rps/MBps. Закладывать:
- CPU на парс/ингест (обычно 15–25% ядёр BE «съедает» активный поток),
- диск под временные файлы и компакшн,
- сеть: ingest-сегмент (25 GbE хорошо).
Оценка ingest-мощности (грубо):
- NVMe-узел с 32 vCPU/256 GB RAM при RF=3 устойчиво переваривает 200–400k событий/сек (плоских JSON/CSV) при умеренных трансформациях и разумном партиционировании.
- При PK-таблицах (upsert) и высокой кардинальности скорость ниже: учитывайте -30–50% от плоских вставок.
5) FE sizing (не экономим на «мозгах»)
- 3 FE: по 8 vCPU / 16–32 GB RAM, локальный SSD для журналов.
- Heap JVM под FE: стартовать 8–16 GB, следить за GC паузами; Metastore/каталоги/большое количество таблиц/MV требуют памяти.
- Если таблиц/партиций/таблетов очень много (сотни тысяч), повышайте RAM FE и ускоряйте SSD.
6) Параметры данных, влияющие на сайзинг
-
Тип таблицы:
- Duplicate Key — быстрые append, дёшево по CPU, но растит объём.
- Aggregate Key — экономит место/время запроса, но усложняет ingest/логическую консистентность.
- Primary Key — удобный upsert/delete, но дороже по CPU/диску (индексы, слияния).
- Партиционирование: «дата-плюс» (день/час) для RT/инкремента; не плодите тысячи партиций без нужды.
- Дистрибуция (hash по бизнес-ключу) — балансирует нагрузку; избегайте skew-ключей.
- Материализованные представления: дают мощный выигрыш по CPU/латентности, но сто́ят диска и времени обновления — учитывайте их как отдельные наборы.
7) Сборка референс-профилей
Small (POC / отдел BI, ≤10 TB сжатых)
- FE: 3 × (4–8 vCPU, 16 GB RAM, SSD 100 GB)
- BE: 3 × (16 vCPU, 64–128 GB RAM, NVMe 2×1.92 TB, 10 GbE)
- RF=3, TTL на сырые данные 90 дней, 1–2 MV на витрину.
Medium (корпоративный BI, 10–80 TB сжатых, 100–300 активных пользователей)
- FE: 3–4 × (8 vCPU, 32 GB RAM, SSD ≥ 200 GB)
- BE: 6–12 × (32 vCPU, 256 GB RAM, NVMe 2–4×3.84 TB, 25 GbE)
- RF=3, выделенная сеть под ingest, compaction-окна ночами, 10–30 MV.
Large (RT аналитика + DWH-витрины, 80–300 TB сжатых, высокий ingest)
- FE: 4–5 × (16 vCPU, 64 GB RAM, быстрые SSD)
- BE: 16–40 × (32–64 vCPU, 256–512 GB RAM, NVMe 4–8×3.84–7.68 TB, 25–40 GbE)
- RF=3, много MV, частичный tiering в объектное хранилище, агрессивный TTL.
8) Пример расчета (workload-driven)
Вводные:
- Исторические сырые данные Raw_hist = 120 TB
- Прирост Raw_new = 1.2 TB/день
- Горизонт планирования = 365 дней
- Compression = 3.5, RF=3, Overhead=0.15, Headroom=0.3
- MV суммарно добавят ~30% к сжатому объему активной зоны (по ключевым витринам)
- BI: 200 одновременных пользователей, SLA P95 ≤ 3 сек для типовых дашбордов
- RT ingest из Kafka: 250k событий/сек в часы пик (плоские записи)
Хранилище:
Raw_total = 120 + 1.2*365 = 120 + 438 = 558 TB Compressed = 558 / 3.5 ≈ 159.4 TB + Overhead (15%) → 183.3 TB + MV (добавим 30% к активным 40 TB, допустим активная зона 40 TB → 12 TB) Итого до RF ≈ 195.3 TB RF=3 → 585.9 TB Headroom 30% → 761.7 TB Округлим до ~780 TB брутто на кластер.
BE-узлы (NVMe):
Берем узлы по 4×7.68 TB NVMe ≈ 30.7 TB брутто → с учётом файловой системы/резерва ~25 TB полезной.
Нужно 780 / 25 ≈ 31–32 BE только под хранение.
С учетом CPU/ingest/запаса возьмем 36 BE (32 vCPU, 256 GB RAM, 25 GbE).
FE-узлы:
4 FE по 16 vCPU / 64 GB RAM, быстрые SSD, LB перед ними.
Пропускная ingest-способность:
36 BE такого класса с запасом выдержат 250k ev/s (при RF=3), за счет распределения по партициям/таблетам. Закладываем Routine Load с batch 20–50k, контролируем compaction-окна.
Конкуренция запросов:
200 одновременных BI-пользователей → активных тяжелых ~30–40.
На 36 BE с 32 vCPU (итого 1152 vCPU) можно держать 40 тяжёлых + сотни легких, сохранив SLA при наличии MV.
9) Кубики настройки, влияющие на емкость/производительность
- Партиции: день/час под ingest; не превышать «тысячи» партиций без нужды.
- Tablets per BE: стремиться к ровному распределению; если перекос — ребаланс.
- Compaction: планировать окна (ночные/низкая нагрузка), следить за очередями, не давать автокомпакшну «душить» ingest.
- Materialized Views Rewrite: проверять, что запросы реально переписываются в MV (EXPLAIN/PROFILE).
- Session vars: лимитировать ресурсы тяжёлым пользователям (query_timeout, mem_limit, parallel_fragment_exec_instance_num).
- Kafka: включать backpressure/лимиты, использовать несколько consumer-групп для параллелизма.
10) Риски и как застраховаться
-
Скью по ключу дистрибуции → «горячие» BE.
Меры: анализ распределения ключей, выбор лучшего hash-ключа, переразбиение. -
Взрыв партиций/таблетов → рост метаданных/FE-нагрузки.
Меры: укрупнить партиции, архивировать старые, слить мелкие tablets. -
Перекрытие compaction и ingest → лаги в загрузке.
Меры: расписание компакшна, лимиты, ночные окна, больше дисков/BE. -
MV стареют/не обновляются вовремя → несвежие дашборды.
Меры: инкрементальные MV, мониторинг лагов, явные триггеры обновлений. -
FE узкие места (GC/метаданные) → рост планирования.
Меры: больше RAM/heap, быстрые SSD, горизонтальный масштаб, чистка «мусорных» объектов. -
Недооценка RF и headroom → out-of-space при пиках и ребалансе.
Меры: RF=3 в проде, 30% свободного, capacity-alerts.
11) Kubernetes/облако (коротко)
- FE — StatefulSet с anti-affinity, диски SSD для журнала.
- BE — DaemonSet/StatefulSet с локальными NVMe, topology-spread, pod-disruption-budgets.
- Сетевые классы 25 GbE эквивалент, node-pools под ingest и под BI.
- В облаках: выбирайте инстансы с «настоящим» NVMe (или предоставляемыми высокими IOPS), отдельные группы под FE и BE.
12) Чек-лист перед финализацией сайзинга
- Согласован ретеншн/TTL и доля агрегированных витрин.
- Подтверждена компрессия на пилоте (прогон 2–3 типовых таблиц).
- Протестирован ingest в пике (Kafka/Batch) на 48–72 часа.
- Замерена конкуренция типовых дашбордов, зафиксированы SLA.
- Внесены MV и проверен rewrite.
- Заложен RF, headroom 30%, окна компакшна.
- Мониторинг FE/BE/compaction/lag готов в день 0.




