StarRocks architecture
StarRocks is a distributed analytical database designed for distributed analytical workloads and high-performance SQL query processing. It uses a massively parallel processing (MPP) architecture, where a query is divided into multiple execution tasks and processed concurrently by multiple compute nodes.
In ADH, StarRocks is deployed using the shared-data architecture. This architecture separates persistent data storage from query processing: data is stored in an external storage, while StarRocks compute nodes process queries.
|
IMPORTANT
The shared-nothing model and Backend (BE) nodes are not available in the ADH deployment of StarRocks. ADH uses only the shared-data model with FE and CN components.
|
An ADH StarRocks cluster consists of the following components:
-
Frontend (FE) — manages metadata, accepts client connections, creates and optimizes query execution plans, and coordinates query execution.
-
Compute Node (CN) — executes query fragments and provides local data caching.
Key features
This approach provides the following capabilities:
-
Massively parallel query execution. Query workloads are distributed across multiple CN nodes and processed concurrently.
-
Fully vectorized execution engine. Data is being stored and processed in a columnar manner, which makes the usage of CPU processing power more efficient.
-
Elastic compute scaling. CN nodes can be added or removed without redistributing data between compute nodes.
-
Local data caching. Frequently accessed data can be cached on CN nodes to reduce access to remote storage.
-
SQL interface. StarRocks supports the MySQL protocol, allowing standard SQL clients and analytical tools to connect to the cluster.
-
Query optimization. FE nodes use Cost-Based Optimizer (CBO) to check SQL statements and generate execution plans optimized for distributed processing.
-
High availability. Multiple FE nodes can be deployed to provide metadata redundancy and service continuity.
-
Resource isolation. Storage and compute resources can be managed independently, which is particularly useful for analytical workloads with variable compute requirements.
Architecture overview
StarRocks represents the control and compute layers in the logical scheme of an SQL query:
-
Client layer. SQL clients, BI applications, and other applications that submit queries to StarRocks.
-
Control layer. FE components that manage metadata, plan queries.
-
Compute layer. CN components that execute workloads.
-
Storage layer. External storage that contains persistent StarRocks data.
Frontend (FE)
The Frontend (FE) is the control-plane component of StarRocks. It coordinates the cluster and manages the metadata required to process SQL requests.
The main FE responsibilities include:
-
Managing cluster and database metadata.
-
Accepting client connections.
-
Managing client sessions.
-
Parsing SQL statements.
-
Creating logical query plans.
-
Optimizing queries.
-
Generating physical execution plans.
-
Scheduling query execution.
-
Coordinating distributed query processing.
Each FE maintains a copy of the StarRocks metadata locally. Multiple FE nodes can be deployed to improve availability and distribute workloads.
FE nodes can operate in different roles.
-
Leader — performs metadata write operations and coordinates metadata changes.
-
Follower — maintains a synchronized copy of metadata and participates in leader election.
-
Observer — maintains metadata and can increase the query-serving capacity of the FE layer but does not participate in leader election.
If the current leader becomes unavailable, eligible follower FEs can elect a new leader using the Raft-based metadata replication mechanism.
From the perspective of query processing, the FE is responsible for deciding how a query should be executed. It does not perform the main data-processing workload itself. Instead, it creates an execution plan and distributes the corresponding work to CN nodes.
Compute Node (CN)
Compute Node (CN) is the main execution component of the StarRocks service. CN nodes are responsible for executing the physical query plan generated by the FE. They perform operations such as:
-
Scanning data.
-
Filtering records.
-
Joining datasets.
-
Aggregating results.
-
Sorting data.
-
Executing expressions and functions.
-
Exchanging intermediate data with other CN nodes.
-
Caching frequently accessed data.
CN nodes are designed for shared-data deployments, they do not store persistent data.
External storage
In ADH, StarRocks deployments support HDFS and Ozone as external storages. If both HDFS and Ozone are installed, HDFS is selected as the external storage.
StarRocks can be configured to use other object storage options, such as S3-compatible systems, GCS, or Azure.
Component interaction
Query processing
FE and CN components work together during query processing.
A typical query lifecycle consists of the following stages:
-
Client connection. A client application sends an SQL statement to an FE node. The FE becomes the entry point for the query and coordinates subsequent processing.
-
SQL parsing and analysis. The FE parses the SQL statement and validates the requested operation against the available metadata. It determines referenced table and columns, analyzes filtering conditions, and identifies the required aggregation.
-
Query optimization. The FE creates and checks the query plan. The optimizer determines how the requested operations can be executed efficiently across the CN cluster. The resulting physical plan consists of execution fragments that can be distributed between CN nodes.
-
Query scheduling. The FE schedules the execution fragments on available CN nodes.
-
Data access. When a CN starts executing a query fragment, it accesses the required data. The data may already be available in the local cache. If it is not, the CN retrieves the required data from external storage.
-
Distributed execution. CN nodes execute their assigned query fragments in parallel. Individual CNs can process different portions of the query simultaneously.
-
Intermediate data exchange. Some operations require CN nodes to exchange intermediate results. This is common for distributed joins and aggregations.
-
Result delivery. After the required execution fragments have completed, the FE coordinates the query result returned to the client.
Data ingestion
FE and CN nodes also cooperate when data is loaded into StarRocks:
-
The FE receives and coordinates the ingestion request. It determines the appropriate execution strategy and assigns the workload to CN nodes.
-
CN nodes perform the data-processing part of the ingestion operation and write the resulting persistent data to the external storage layer.
Scaling
One of the main advantages of the shared-data architecture is the ability to scale the compute layer independently.
When query concurrency or computational requirements increase, additional CN nodes can be added to a cluster. Because persistent data is stored externally, adding a CN does not require moving existing data between compute nodes.
This separation is particularly useful when storage requirements and query-processing requirements grow at different rates.
High availability
High availability of StarRocks can be provided in all functional layers:
-
FE availability. Multiple FE nodes can be deployed to maintain metadata availability. Observer nodes replicate metadata, and follower nodes participate in leader election. If the active leader fails, another eligible follower can become the new leader.
-
CN availability. In case a CN becomes unavailable, another available CN can process subsequent workloads. CN availability affects the compute capacity of a cluster. Therefore, the number of CN nodes should be planned according to expected query concurrency and workload requirements.
-
Storage availability. Persistent data is maintained by an external storage system. The availability, durability, capacity, and performance of that storage system are part of the overall StarRocks service characteristics.