Solr Cluster Types
A Solr cluster is a group of servers that each run one or more Solr nodes.
There are two general modes of operating a cluster of Solr nodes. One mode provides central coordination of the Solr nodes (SolrCloud Clusters), while the other allows you to operate a cluster without this central coordination (User-Managed Cluster).
| "User Managed" is sometimes referred to as "Standalone" in source code. |
Both modes share general concepts, but ultimately differ in how those concepts are reflected in functionality and features.
First let’s cover a few general concepts and then outline the differences between the two modes.
Cluster Concepts
Servers and Nodes
A server is the hardware or virtual machine that hosts Solr software. A node is an instance of a running Solr process that services search and indexing requests. In special cases where oversized pre-existing hardware must be utilized, a server might host two or more nodes. Note that such configurations are typically sub-optimal.
Shards
In both cluster modes, shards are logical divisions of a collection of documents. Shards slice a collection of documents into discrete non-overlapping subsets, and may be based on data values you specify or ranges of a hash on the document ID.
The number of shards determines the theoretical limit to the number of documents that can be stored. It also dictates the amount of parallelization possible for an individual search request.
Replicas
A shard is a logical concept—a slice of your collection. A replica is the physical manifestation of that logical shard. It’s comprised of a Lucene index on disk, and supporting infrastructure for managing the index. A Solr node hosts each of its replicas that are solely responsible for operating that index including handling requests to it for indexing & searching.
A shard must have at least one replica to exist physically. If you have one shard with one physical copy, you have one replica. If you add redundancy by creating additional copies of that shard, you have multiple replicas—each is equally a replica, including the first one.
| There is no "original shard" separate from its replicas. The replicas ARE how the shard exists. This is why we say "a shard with 2 replicas" has 2 total physical copies, not an original plus 2 additional copies. |
All replicas of the same shard contain the same subset of documents. All replicas of the same collection use the same configuration and configset.
The number of replicas determines the level of fault tolerance the cluster has in the event of a node failure. It also dictates the theoretical limit on the number of concurrent search requests that can be processed under heavy load.
Leaders and Followers
Among the replicas for a given shard, one replica is designated as the leader. The leader serves as the source-of-truth for its shard. When document updates are made, they are first processed by the leader replica and then propagated to the other replicas (the exact mechanism varies by cluster mode).
The replicas which are not leaders are called followers.
Cores
In Solr’s implementation, each replica is represented as a core. The term "core" is primarily an internal implementation detail—when you create a replica, Solr creates a core to represent it. Multiple cores can be hosted on any one node.
| Historically the term "core" has mostly been used as a synonym for replica, but the term "core" can be confusing because in everyday English it implies something central and singular. Since there may be many replicas in Solr, and they are distributed across the cluster "Replica" is the preferred term. Core is mostly only used for historical reasons in the code base and other places where renaming things would be disruptive. |
Collections and Indexes
A collection is the complete logical set of searchable documents that share a schema and configuration. In SolrCloud mode (described below), a collection encompasses all the shards and their replicas.
An index refers to the physical data structures written to disk by Apache Lucene. Each replica maintains exactly one Lucene index on disk, containing the actual inverted indexes, stored fields, and other data structures that enable search.
This creates a clear hierarchy from logical concepts to physical storage:
Cluster
└─> Collection (logical grouping of all searchable documents)
└─> Shard 1 (logical partition)
│ └─> Replica 1 / Core 1 (physical instance)
│ │ └─> Lucene Index (disk structures)
│ └─> Replica 2 / Core 2 (physical instance)
│ └─> Lucene Index (disk structures)
└─> Shard 2 (logical partition)
└─> Replica 1 / Core 3 (physical instance)
│ └─> Lucene Index (disk structures)
└─> Replica 2 / Core 4 (physical instance)
└─> Lucene Index (disk structures)
In this example, a collection is divided into 2 shards, each shard has 2 replicas for redundancy, and each replica maintains its own Lucene index on disk.
SolrCloud Clusters
A SolrCloud cluster (or simply "SolrCloud") uses Apache ZooKeeper to provide the centralized cluster management that is its main feature. ZooKeeper holds a list of each live node in the cluster and the state of each collection (and thus shard and replica).
In this mode, configuration files (a "Configset") are stored in ZooKeeper and not on the file system of each node. Like most things in Solr, this choice is configurable. When configuration changes are made, they must be uploaded to ZooKeeper, which in turn makes sure each node knows changes have been made.
SolrCloud manages collections as first-class entities.
A collection represents the entire group of shards and replicas that together provide access to a corpus of documents.
Collections share the same configurations (schema, solrconfig.xml, etc.).
This centralization of cluster management means that operations can be performed on the entire collection at one time.
When changes are made to configurations, a single command to "reload" the collection will automatically reload each individual core (replica) that is a member of the collection with the latest configuration.
Collections may also be configured to provide automatic routing of documents to shards by hashing document ids and automatically assigning ranges of the possible hash values to shards. Some degree of control over what documents are stored in which shards is also available, if needed.
Incoming requests, either to index documents or for user queries, can be sent to any node of the cluster and Solr will consult the state information in ZooKeeper in order to route the request to an appropriate replica of each shard.
In SolrCloud, the leader replica within a shard is flexible, with built-in mechanisms for automatic leader election in case the current leader fails. This means another replica can become the leader, and from that point forward it is the source-of-truth for all other replicas of that shard.
As long as one replica of each relevant shard is available, a user query or indexing request can still be satisfied when running in SolrCloud mode.
User-Managed Cluster
User-managed mode has no concept of a collection as a managed entity, so for all intents and purposes each Solr core is configured and managed independently. Only configuration parameters keep related cores from behaving as completely independent entities.
Solr’s user-managed cluster is nothing more than a set of Solr nodes in standalone mode, and thus know nothing of collections or shards, and Zookeeper is not used as a centralized storage for any configuration or real-time state. Instead the "user" / operator (you) must do-it-yourself, both at client layers and probably some local scripts to automate whatever needs doing. The Core Admin APIs of Solr with some special core-level APIs like replication become the fundamental building blocks for you to design / build a system to your own specifications. This was what people did prior to SolrCloud, and some still do.
If the corpus of documents is too large for a single shard, the logic to create multiple shards is entirely left to the user. There are no automated or programmatic ways for Solr to create shards on demand.
Routing documents to shards is handled manually, either with a hashing system (that you design and implement), assignment of documents to shards based on the value of a field (implicit routing), or a simple round-robin list of shards that sends each document to a different shard. Document updates must be sent to the right shard or duplicate documents could result.
In user-managed mode, the distinction between leader and follower replicas becomes critical. Identifying which node will host the leader replica and which host(s) will have follower replicas dictates how each node is configured. In this mode, all document updates are sent to the leader replica only. Once the leader has completed indexing, each follower replica will request the index updates and copy them from the leader.
Load balancing is achieved with an external tool or process, unless request traffic can be managed by the leader or one of its follower replicas alone.
If the leader replica goes down, there is no built-in failover mechanism.
A follower replica could continue to serve queries if the queries were specifically directed to it.
Promoting a follower replica to serve as the leader would require changing solrconfig.xml configurations on all replicas and reloading each core.