Introduction: Neo4j clustering architectureEnterprise Edition
Overview
Neo4j’s clustering provides these main features:
-
Scalability: A Neo4j cluster is a set of servers running multiple databases. Servers and databases are decoupled: servers provide computation and storage power for databases to use. Each database has its own independent topology, organized into primaries and secondaries (for read scaling).
-
Fault tolerance: Primary database allocations provide a fault tolerant platform for transaction processing. A database remains available for writes as long as a simple majority of its primary allocations are functioning.
-
Operability: Database management is separated from server management. For details, see Managing databases in a cluster and Managing servers in a cluster.
-
Causal consistency: When invoked, a client application is guaranteed to read at least its own writes.
For information about cluster design patterns and anti-patterns, see Designing resilient multi-region cluster deployments.
Operational view
From an operational point of view, it is useful to view the cluster as a homogenous pool of servers which run a number of databases.
Note that primary and secondary are roles for a copy of a database.
Servers do not have roles.
Instead, they can be constrained by a modeConstraint set to PRIMARY, SECONDARY, or NONE.
That means they can host copies of standard databases that have either a primary role or a secondary role (see servers 1 or 3 in the Figure 1).
If the modeConstraint is set to NONE, a server can host copies of databases that have either role.
In other words, a server can host primaries for some databases and secondaries for other databases (see server 6 in the Figure 1).
Similarly, it is possible for a database to be hosted on only one server, even when that server is part of a cluster. In such cases, the database allocation is always primary. See single primary for details.
Database primaries
A primary is a copy of a database that is a participant in the processing of write operations and can be the writer for that database. A database can have one or more primary allocations within a cluster.
Only primaries are eligible to act as the writer for a database. At any given time, only one primary is automatically elected as a writer among the database’s primaries. The writer may change over time.
The database writer synchronously pushes writes to other primaries and does not allow a commit to be completed until it receives confirmation that the data has been written to enough members.
For high availability, create a database with multiple primaries. If high availability is not required, a database can be created with a single primary to achieve minimum write latency.
If too many primaries fail, the database can no longer process writes and becomes read-only.
How many primaries you should have
Primaries provide:
-
The ability to write to your database (with optional fault tolerance).
-
The ability to read from your database.
-
Fault tolerance for different failure scenarios.
The fault tolerance is calculated with the formula:
M = 2F + 1, where M is the number of primaries required to tolerate F faults.
Generally speaking, fault tolerance is the number of primary database copies you can lose without affecting a certain operation. For instance, with one primary copy, you have no fault tolerance, because if it goes offline, nothing is available. If you have two servers, each with a primary copy of the database, and one goes offline, the other will still have some copy of the data, so read availability would be preserved.
Types of fault tolerance, listed from easiest to hardest to lose, are as follows:
-
Write availability:
If write availability is lost, your database cannot accept any more writes. -
Read availability:
If read availability is lost, your database cannot serve any more reads. -
Durability:
If durability is lost, the data written to your database is lost, and you need to restore the database from a backup.
Operations such as shutting Neo4j process down to upgrade the binaries, or taking the server it runs on offline for maintenance, are included as faults in this context, since they make the database copy unavailable.
If you want upgrades with no downtime, you need fault tolerance.
Therefore, you need minimum three primaries (and three servers to host each copy) to be able to maintain write availability with the failure of one member. This is enough for most deployments.
If you want to retain write availability with the failure of two primary members, you need five primaries.
The maximum number of primaries you can have is 11, but it is not recommended having that many primaries. Because the more primaries you have, the more servers you have to contact for each write operation, which can increase the latency of writes.
Database secondaries
A secondary is a database copy asynchronously replicated from primaries via transaction log shipping. Secondaries periodically check an upstream database member for new transactions, which are then transferred to them.
The main purpose of database secondaries is to scale out read workloads. Secondaries act like caches for graph data and can execute arbitrary read-only queries and procedures.
Multiple secondaries can be fed data from a relatively small number of primaries, providing a significant distribution of query workloads for better scalability. Databases can have a fairly large number of secondaries.
The loss of a secondary does not affect the database’s availability; however, it reduces the query throughput. It also does not affect the database fault tolerance.
While secondaries serve as a copy of your database, providing some level of durability (what is committed cannot be lost), they do not guarantee it completely. Secondaries pull updates from a selected upstream member on their own schedule. Due to their asynchronous nature, secondaries may temporarily lag behind the primary, meaning recently committed transactions may not be immediately visible on secondaries.
How many secondaries you should have
Secondaries typically provide read scaling, i.e. if you have more read queries happening than your primaries can handle, you can add secondaries to share the load; or even configure the query routing so that reads preferentially target secondaries to leave the primaries free to handle just the write workload.
So, there is no hard rule about the number of secondaries to have. Starting with zero and adding more until your read performance and cluster stability is acceptable is the usual approach.
The maximum number of secondaries you can have is 20.
Primaries and secondaries for the system database
The system database, which records what databases are present in the DBMS, also can be in a primary or secondary mode.
However, unlike standard databases, it is not configured using Cypher commands to define the topology.
Instead, it is controlled through the server.cluster.system_database_mode setting.
Use the following guidelines when deciding how many primary and secondary system databases to have and which servers should host them:
-
Stable, long-lived servers are good candidates to host a
systemprimary, since they are expected to remain online and can be intentionally shut down when needed. -
Ephemeral or frequently changing servers are good candidates to host a
systemsecondary, as they may be added or removed more often. -
A single
systemprimary provides no fault tolerance for writes to thesystemdatabase. Therefore, in a typical cluster deployment, it is best to start with threesystemprimaries to ensure write availability. -
Although the write volume for the
systemdatabase is low and it can tolerate higher write latency, allowing to have more than 11systemprimaries, doing so is generally not recommended.
Examples of database topologies
For information about the cluster deployment across multiple data centers, refer to the Designing resilient multi-region cluster deployments.
Single primary
If you have a single copy of the database, it is a primary. All writes and reads goes through this copy. If the copy becomes unavailable, no writes or reads are possible. If the disk for that copy is lost or corrupted, durability is lost and you must restore from the latest full backup.
Three primaries
In a cluster with three primaries, the members elect a leader to process write operations. Each write is replicated to at least one additional primary before being considered committed (durable). This ensures that if any single primary fails, that update remains available on another member. That includes if the database copy is fully lost, including the disk being unrecoverable.
If one primary copy fails, the database is still write-available with the remaining two primaries, but it no longer has fault tolerance for its write availability.
Another failure would prevent any new writes from being processed until either one of the other members is brought back, or the database is recreated with new members. The database would still be read-available on the last member though.
The non-writer primaries also provide read capacity and fault tolerance.
By default, read queries are routed away from the writer (see dbms.routing.reads_on_writers_enabled).
See Geo-distribution of user database primaries for the pattern of deploying a cluster with three primaries across three data centers.
Five primaries
If you want fault tolerance greater than one arbitrary database member, deploy five primaries. You will have tolerance to the failure of any two primaries. The remaining three primaries can still maintain quorum and ensure the database continues to operate.
For information about deploying five primaries across multiple data centers, see Resilient multi-region cluster deployment → Designing a resilient multi-data center cluster.
Single primary plus secondaries
As described above, a single primary provides no fault tolerance for both write availability or durability. If the single primary fails, no write operations can be processed, and if its disk is lost, the most recent updates may be lost.
However, adding one or more secondaries means that read availability can be maintained despite the loss of the primary.
Keep in mind that the secondaries may not have the most up to date data, which is only guaranteed to be present on the primary.
The secondaries typically handle all of the read queries (see dbms.routing.reads_on_writers_enabled).
Three primaries plus secondaries
As described above, three primaries provide fault tolerance for both write availability and durability.
Both secondaries and non-writer primaries can handle read queries.
Although, you can configure only secondaries to handle reads (see dbms.routing.reads_on_primaries_enabled) if you need primaries to focus on the write workload.
The loss of any single database copy does not affect write availability, read availability, or durability.
If all the secondaries fail, then primaries that are not acting as the writer start handling read queries to maintain read availability.
Causal consistency
While the operational mechanics of the cluster are interesting from an application point of view, it is also helpful to think about how applications use the database to get their work done. In many applications, it is typically desirable to both read from the graph and write to the graph. Depending on the nature of the workload, it is common to want reads from the graph to take into account previous writes to ensure causal consistency.
|
Causal consistency is one of numerous consistency models used in distributed computing. It ensures that causally related operations are seen by every instance in the system in the same order. Consequently, client applications are guaranteed to read their own writes, regardless of which instance they communicate with. This simplifies interaction with large clusters, allowing clients to treat them as a single (logical) server. |
Causal consistency makes it possible to write to databases hosted on servers in primary mode and read those writes from databases hosted on servers in secondary mode (where graph operations are scaled out). For example, causal consistency guarantees that the write which created a user account is present when that same user subsequently attempts to log in.
On executing a transaction, the client can ask for a bookmark which it then presents as a parameter to subsequent transactions. Using that bookmark, the cluster can ensure that only servers which have processed the client’s bookmarked transaction will run its next transaction. This provides a causal chain which ensures correct read-after-write semantics from the client’s point of view.
Aside from the bookmark everything else is handled by the cluster. The database drivers work with the cluster topology manager to choose the most appropriate servers to route queries to. For instance, routing reads to database secondaries and writes to database primaries.
Glossary
- allocator
-
A component in the cluster that allocates databases to servers according to the topology constraints specified and an allocation strategy.
- asynchronous replication
-
Asynchronous replication is used by secondary copies to poll for new transactions, which means they cannot be guaranteed to have received the most recent transactions. This enables efficient scale-out of read-performance.
- Aura instance
-
A fully-managed DBMS represented by a single instance ID, that is running in the Neo4j Aura cloud.
- auto-commit transaction
-
An automatically committed transaction that contains a single query.
- Bolt protocol
-
Bolt is a protocol used for interaction between Neo4j instances and drivers.
- bookmark
-
A marker the client can request from the cluster to ensure that it is able to read its own writes so that the application’s state is consistent and only databases that have a copy of the bookmark are permitted to respond.
- category (Bloom)
-
A category is based on a node label and is defined in a Perspective as a way of visually distinguishing nodes with the same label(s).
- causal consistency
-
All servers in a cluster agree on the order in which transactions take place. The position of a server on the causal chain can be guaranteed using a bookmark.
- cluster
-
A Neo4j DBMS that spans multiple servers working together to increase fault tolerance and/or read scalability. Databases on a cluster may be configured to replicate across servers in the cluster thus achieving read scalability or high availability.
- client application
-
Software that interacts with a Neo4j server.
- commit
-
A commit is the successful completion of a transaction, which ensures durability of any changes made. For more details, visit Operations Manual → Transaction management.
- composite database
-
Composite databases are the means to access partitioned graph data with a single Cypher query.
- constraint
-
Constraints are sets of data modeling rules that ensure the data is consistent and reliable.
- Cypher®
-
Neo4j’s graph query language.
- data model
-
A data model defines how information is organized in a database. A good data model will make querying and understanding your data easier. In Neo4j, the data models have a graph structure.
- database
-
A database is a container used by the DBMS to manage and store graph data. The physical structure of data is controlled by the database.
- database vs graph
-
Databases are the physical containers of graph data. Graphs are the logical structure of data in Neo4j.
- Database Management System
-
Database Management System, or DBMS, capable of managing multiple databases. A DBMS may run on a single server, or span several servers configured as a cluster.
- database schema
-
The prescribed property existence and datatypes for nodes and relationships.
- deallocate
-
An act of removing a database from a server or a server from a cluster without loss of data or reduced fault tolerance.
- degree (of a node)
-
The number of relationships of a specific node; loops are counted twice.
- disaster recovery
-
A manual intervention to restore availability of a cluster, or databases within a cluster.
- driver
-
A software library that provides access to Neo4j from a particular programming language.
- election
-
In the event that the Raft leader becomes unresponsive, followers automatically trigger an election and vote for a new leader.
- entity
-
A node or a relationship.
- expression (Cypher)
-
A component of a Cypher query which produces values. It may be used in projections, as a predicate, or when setting properties on graph elements.
- fabric
-
Fabric is the architectural design of a unified system that provides a single access point to local or distributed graph data.
- fault tolerance
-
A guarantee that a cluster can maintain a database’s persistence and availability in the event of one or more servers failing.
- follower
-
A primary copy of a database acting as a follower, receives and acknowledges synchronous writes from the leader.
- Generative AI (GenAI)
-
A type of artificial intelligence (AI) system that generates text, images, or other media in response to prompts.
- graph
-
A logical representation of a set of nodes where some pairs are connected by relationships.
- index
-
Data structure that improves read performance of a database.
- knowledge graph
-
A specific type of graph that has an organizing principle so that a user (or a computer system) can reason about the underlying data. The organizing principle provides an additional layer of structure that adds context to support knowledge discovery.
- label
-
Marks a node as a member of a named and indexed subset. A node may be assigned zero or more labels.
- leader
-
A single primary copy of a database is designated as the leader. It receives all write transactions from clients and replicates writes synchronously to followers and asynchronously to secondary copies of the database.
- main database
-
In terms of Neo4j Enterprise Studio, the database(s) containing the user’s data. Can exist in the same Neo4j deployment as the tool asset database.
- motif
-
A description of a specific pattern within a graph.
- node
-
A node represents an entity or discrete object in your graph data model. Nodes can be connected by relationships, hold data in properties, and are classified by labels.
- operator
-
A symbol representing a mathematical or logical operation.
- parameter
-
Named value provided when running a Cypher statement.
- path
-
A sequence of nodes and the relationships connecting them, that does not contain duplicate relationships. Several paths can match a pattern.
- pattern
-
A specific arrangement of nodes and relationships that can be matched in a graph. A pattern follows a motif.
- perspective (Bloom)
-
A Perspective defines a certain business view or domain that can be found in the target Neo4j graph. A single Neo4j graph can be viewed through different Perspectives, each tailored for a different business purpose.
- primary
-
A copy of the database that is able to process write transactions and is eligible to be elected as a leader. It participates in fault tolerant writes as it is part of the majority required to acknowledge and commit write transactions.
- primary vs secondary
-
In a cluster, databases can operate in either primary or secondary mode. Primary databases are able to process write and read transactions, ensuring fault tolerance. Secondary databases are replicated asynchronously from primaries, and their main purpose is to provide read scaling within the cluster.
- project (Aura)
-
An isolated environment in the unified Aura console that contains its own database instances, configurations, and resources. Preceded by tenant in the classic Aura console.
- property
-
Properties are key-value pairs that are used for storing data on nodes and relationships.
- query (Cypher)
-
A statement that retrieves or writes information to a database.
- Raft group
-
A group of servers that are participating in hosting a particular database in primary mode.
- Raft group member
-
A server that is participating in a Raft group. A server can be a member of one or more groups.
- Raft log
-
A shared log between all Raft group members that is guaranteed to be consistently updated and viewed by those members. The log contains both database data and operational state of the Raft group.
- Raft protocol
-
The networking mechanism that enables a database to replicate its data across multiple servers to give high availability for accessing the data and high durability to the data stored.
- read scaling
-
Distributing query load by creating additional database copies hosted in secondary mode (read-only).
- relationship
-
A relationship represents a connection between nodes in your graph data model. Relationships connect a source node to a target node, hold data in properties, and are classified by type.
- secondary
-
An asynchronously replicated copy of the database that provides read scaling within the cluster.
- seed
-
A seed is a database dump or a full backup used to create a database on a cluster. This is sometimes called seeding.
- server
-
A physical machine, a virtual machine, or a container running an instance of Neo4j. Servers can be standalone or part of a cluster.
- session
-
A causally linked sequence of transactions.
- session consistency
-
An alternative name for Neo4j’s causal consistency.
- standalone
-
A single server running Neo4j and not part of a cluster.
- synchronous replication
-
Synchronous replication requires the leader primary to replicate a transaction and block the commit until a quorum of the follower primaries acknowledges that the transaction is successfully replicated. Once the transaction is replicated, the commit is allowed to proceed. This ensures data durability and consistency within the cluster.
- system database
-
A database used by Neo4j to store system information.
- tenant (Aura)
-
An isolated environment in the classic Aura console that contains its own database instances, configurations, and resources. Replaced by project in the unified Aura console.
- tool asset database
-
In terms of Neo4j Enterprise Studio, the database where tools' assets are stored. This can be in the same Neo4j deployment as the main database(s) or in a separate deployment.
- topology
-
A configuration that describes how the copies of a database should be spread across the servers in a cluster, see primary mode and secondary mode.
- transaction
-
A transaction comprises a unit of work performed against a database. It is treated in a coherent and reliable way, independent of other transactions. Transactions comply with the ACID consistency model (atomic, consistent, isolated, and durable).