Архитектура StarRocks

StarRocks — это распределенная аналитическая база данных, предназначенная для высокопроизводительной обработки SQL-запросов. Она использует архитектуру массово-параллельной обработки (Massively Parallel Processing, MPP). При этом запрос разбивается на несколько подзапросов, которые одновременно обрабатываются несколькими вычислительными узлами.

В ADH StarRocks развертывается с использованием архитектуры shared-data. Эта архитектура подразумевает разделение обработки запросов и хранения постоянных данных: данные хранятся во внешнем хранилище, а вычислительные узлы StarRocks только обрабатывают запросы.

ВАЖНО
Модель shared-nothing и узлы Backend (BE) недоступны при установке сервиса StarRocks в ADH. В ADH используется только модель shared-data с компонентами FE и CN.

Кластер StarRocks в ADH состоит из следующих компонентов:

  • Frontend (FE) — управляет метаданными, принимает клиентские соединения, создает и оптимизирует планы выполнения запросов и координирует их исполнение.

  • Compute Node (CN) — выполняет фрагменты запросов и может кешировать данные.

Основные возможности

Такой подход предоставляет следующие возможности:

  • Массово-параллельное выполнение запросов. Рабочие нагрузки запросов распределяются между несколькими узлами CN и обрабатываются одновременно.

  • Полностью векторизированный вычислительный движок. Данные хранятся и обрабатываются в колоночном формате, что позволяет использовать процессорные мощности эффективнее.

  • Эластичное масштабирование вычислений. Узлы CN можно добавлять или удалять без перераспределения данных между вычислительными узлами.

  • Локальное кеширование данных. Узлы CN могут кешировать часто используемые данные — это позволяет сократить число обращений к внешнему хранилищу.

  • SQL-интерфейс. StarRocks поддерживает протокол MySQL, по которому стандартные SQL-клиенты и аналитические инструменты могут подключаться к кластеру.

  • Оптимизация запросов. Узлы FE используют анализатор на основе стоимости (Cost-Based Optimizer, CBO) для проверки SQL-инструкций и составления планов выполнения, оптимизированных для распределенной обработки.

  • Высокая доступность. Можно развернуть несколько узлов FE, чтобы обеспечить избыточность метаданных и непрерывность работы сервиса.

  • Изоляция ресурсов. Ресурсами хранения и вычисления можно управлять независимо, что особенно полезно для рабочих нагрузок с меняющимися требованиями к вычислительным ресурсам.

Обзор архитектуры

StarRocks представляет слои управления и вычислений в логической схеме выполнения SQL-запроса:

  • Клиентский слой. SQL-клиенты, BI-приложения и другие приложения, отправляющие запросы в StarRocks.

  • Слой управления. Компоненты FE, управляющие метаданными и формирующие планы запросов.

  • Вычислительный слой. Компоненты CN, выполняющие рабочие нагрузки.

  • Слой хранения. Внешнее хранилище, содержащее постоянные данные StarRocks.

Схема выполнения SQL-запроса с использованием StarRocks
Схема выполнения SQL-запроса с использованием StarRocks
Схема выполнения SQL-запроса с использованием StarRocks
Схема выполнения SQL-запроса с использованием StarRocks

Frontend (FE)

Frontend (FE) — это компонент, представляющий собой уровень управления StarRocks. Он координирует работу кластера и управляет метаданными, необходимыми для обработки SQL-запросов.

Основные задачи FE:

  • Управление метаданными кластера и баз данных.

  • Прием клиентских соединений.

  • Управление сессиями клиентов.

  • Синтаксический разбор SQL-выражений.

  • Создание логических планов запросов.

  • Оптимизация запросов.

  • Формирование физических планов выполнения запросов.

  • Планирование выполнения запросов.

  • Координация распределенной обработки запросов.

Каждый узел FE локально хранит копию метаданных StarRocks. Для повышения доступности и распределения нагрузки можно развернуть несколько узлов FE.

Каждый из узлов FE имеет свою роль.

  • Leader — выполняет операции записи метаданных и координирует изменения метаданных.

  • Follower — поддерживает синхронизированную копию метаданных и участвует в выборе лидера.

  • Observer — хранит копию метаданных лидера и может увеличивать пропускную способность слоя FE по обслуживанию запросов, но не участвует в выборе лидера.

Если текущий лидер становится недоступен, узлы Follower FE могут выбрать нового лидера с использованием механизма репликации метаданных на основе Raft.

С точки зрения обработки запросов FE отвечает за определение способа выполнения запроса. Сам FE не выполняет обработку данных. Вместо этого он создает план выполнения и распределяет соответствующие задачи между узлами CN.

Compute Node (CN)

Compute Node (CN) — основной компонент выполнения сервиса StarRocks. Узлы CN отвечают за исполнение физического плана запроса, созданного FE. Они производят такие операции, как:

  • Сканирование данных.

  • Фильтрация записей.

  • Объединение наборов данных.

  • Агрегация результатов.

  • Сортировка данных.

  • Выполнение выражений и функций.

  • Обмен промежуточными данными с другими узлами CN.

  • Кеширование часто используемых данных.

Узлы CN используются при установке согласно модели shared-data, они не хранят постоянные данные.

Внешнее хранилище

В ADH StarRocks поддерживает HDFS и Ozone в качестве внешних хранилищ данных. Если установлены и HDFS, и Ozone, в качестве внешнего хранилища по умолчанию выбирается HDFS.

StarRocks можно настроить для использования других объектных хранилищ, таких как S3-совместимые системы, GCS или Azure.

Взаимодействие компонентов

Обработка запросов

Компоненты FE и CN вместе выполняют обработку запросов.

Типичный жизненный цикл запроса состоит из следующих этапов:

  1. Подключение клиента. Клиентское приложение отправляет SQL-инструкцию на узел FE. FE становится точкой входа для запроса и координирует последующую обработку.

  2. Синтаксический разбор и анализ SQL. FE разбирает SQL-инструкцию и проверяет запрошенную операцию на основе доступных метаданных. Он определяет используемые таблицы и столбцы, анализирует условия фильтрации и требуемую агрегацию.

  3. Оптимизация запроса. FE создает и проверяет план запроса. Оптимизатор определяет, как запрошенные операции могут быть эффективно выполнены в кластере CN. Полученный физический план состоит из фрагментов выполнения, которые могут быть распределены между узлами CN.

  4. Планирование выполнения запроса. FE составляет план выполнения фрагментов запроса на доступных узлах CN.

  5. Доступ к данным. Когда CN начинает выполнять фрагмент запроса, он получает доступ к необходимым данным. Данные могут уже находиться в локальном кеше. Если данных в кеше нет, CN получает необходимые данные из внешнего хранилища.

  6. Распределенное выполнение. Узлы CN параллельно выполняют порученные им фрагменты запроса. Отдельные CN могут одновременно обрабатывать разные части запроса.

  7. Обмен промежуточными данными. Некоторые операции требуют от узлов CN обмена промежуточными результатами. Это типично для распределенных операций объединения и агрегации.

  8. Возврат результата. После выполнения необходимых фрагментов FE координирует возврат результата выполнения фрагментов запроса клиенту.

Загрузка данных

Узлы FE и CN также взаимодействуют при загрузке данных в StarRocks:

  • FE получает запрос на загрузку данных и координирует его выполнение. Он определяет соответствующую стратегию выполнения и назначает рабочую нагрузку узлам CN.

  • Узлы CN выполняют обработку данных в рамках операции загрузки и записывают полученные постоянные данные во внешний слой хранения.

Масштабирование

Одним из основных преимуществ архитектуры shared-data является возможность независимо масштабировать вычислительный слой.

При увеличении количества параллельных запросов или требований к вычислительным ресурсам в кластер можно добавить дополнительные узлы CN. Поскольку постоянные данные хранятся во внешнем хранилище, добавление CN не требует перемещения существующих данных между вычислительными узлами.

Такое разделение особенно полезно, когда требования к объему хранения и обработке запросов растут с разной скоростью.

Высокая доступность

Высокая доступность StarRocks может быть обеспечена на всех функциональных слоях:

  • Доступность FE. Для обеспечения доступности метаданных можно развернуть несколько узлов FE. Узлы Observer FE реплицируют метаданные, а узлы Follower участвуют в выборе лидера. Если активный лидер выходит из строя, другой подходящий узел Follower может стать новым лидером.

  • Доступность CN. Если узел CN становится недоступным, другой доступный CN может начать обрабатывать последующую рабочую нагрузку. Доступность CN влияет на вычислительную емкость кластера. Поэтому количество узлов CN следует планировать с учетом ожидаемого количества параллельных запросов и требований рабочих нагрузок.

  • Доступность хранилища. Постоянные данные хранятся во внешней системе хранения. Доступность, надежность, емкость и производительность этой системы хранения влияет на общие характеристики сервиса StarRocks.

Нашли ошибку? Выделите текст и нажмите Ctrl+Enter чтобы сообщить о ней