docker swarm distributedsystem
Core Idea
A network partition can leave multiple nodes each believing it is the leader; split brain corrupts state unless quorum rules and heartbeats keep a single winner in charge.
- How it happens in systems like Docker Swarm, and the consequences: data corruption, conflicting operations, duplicate tasks.
- Prevention: majority (quorum) rule and leader election with heartbeats.
- Example scenarios with three manager nodes and with four (even-numbered) managers.
Split Brain in distributed systems, such as Docker Swarm, occurs when a network partition causes nodes to lose communication with one another.
resulting in two or more subsets of nodes thinking they are the leader or primary controller of the cluster. This inconsistency can lead to data corruption, conflicting operations, or duplicate tasks being executed.
How it Happens
- Network Partition: A temporary network failure splits the nodes into two or more isolated groups.
- Leader Election Conflict: Each isolated group might independently attempt to elect a leader because they can no longer detect the original leader.
- Independent Decisions: The groups operate as separate clusters and make decisions without synchronization, leading to inconsistent states
In a Docker Swarm cluster:
- Nodes are divided into managers and workers.
- Managers coordinate and maintain the state of the swarm, including service orchestration.
- If a network partition occurs:
- The managers in each partition may attempt to elect a leader.
- This can result in multiple active leaders (split brain), where each isolated group independently schedules and orchestrates tasks, causing service conflicts.
How etcd Handles Split-Brain Scenarios on K8s
etcd, a distributed key-value store used extensively in Kubernetes and other systems, addresses split-brain scenarios through the Raft consensus algorithm. This algorithm enforces rules that require a majority (quorum) of nodes to agree on any changes, thereby maintaining consistency even in the face of network partitions.
Consequences of Split Brain
- Data Inconsistency: Multiple leaders may cause conflicting updates to shared resources.
- Duplicate Workloads: Services or tasks might be scheduled multiple times on different partitions.
- Unrecoverable State: Resolving split brain can be tricky if critical decisions or updates were made independently by both sides.
- Impact on System Reliability
Prevention
Docker Swarm uses mechanisms like:
- Raft Consensus Algorithm: Ensures only one leader can exist at any time.
- Quorum: A majority of managers must agree on cluster state and leader elections. If quorum is lost, no new leader can be elected, and the cluster pauses operations to avoid split brain.
- Network Redundancy: Ensures high availability and prevents network partitions.
Quorum
To identify which partition has the majority (quorum), distributed systems like Docker Swarm use the Raft Consensus Algorithm. Here’s how it works in the event of a network partition:
1. Majority (Quorum) Rule
- Quorum is achieved when more than half (50% + 1) of the manager nodes are in agreement.
- For example:
- In a 3-manager cluster → quorum = 2 nodes.
- In a 5-manager cluster → quorum = 3 nodes.
When a partition occurs:
- Each partition checks how many manager nodes it can communicate with.
- The partition with quorum (the majority of managers) is the active group and continues functioning as the cluster leader.
- The partition without quorum becomes inactive or read-only.
| Managers | Majority | Fault Tolerance |
|---|---|---|
| 1 | 1 | 0 |
| 2 | 2 | 0 |
| 3 | 2 | 1 |
| 4 | 3 | 1 |
| 5 | 3 | 2 |
| 6 | 4 | 2 |
| 7 | 4 | 3 |
2. Leader Election and Heartbeats
- Managers constantly exchange heartbeat messages to confirm each other’s presence.
- When a network partition happens:
- Managers that cannot hear the leader assume it has failed.
- A new leader election begins in each partition using Raft.
The rules for leader election:
- A manager can only become a leader if it has quorum.
- If a partition cannot form a quorum, it cannot elect a leader and goes inactive.
What is heartbeats ?
In Docker Swarm, a “heartbeat” refers to a periodic communication check between nodes (manager nodes specifically) within a swarm cluster, used to ensure that each node is still active and reachable, effectively acting as a mechanism to detect if a node has gone down and needs to be removed from the cluster or replaced; this communication is vital for maintaining the overall health and consistency of the Swarm cluster
Example Scenarios
Scenario: 3 Manager Nodes
- If there is a network partition:
- Partition A: 2 nodes (majority → quorum met).
- Partition B: 1 node (minority → no quorum).
- Partition A will continue to operate (read-write), while Partition B will remain read-only.
Scenario: 4 Manager Nodes (Even Number)
- If split into two partitions:
- Partition A: 2 nodes.
- Partition B: 2 nodes.
- Neither has a majority (quorum = 3). Both groups fail to elect a leader, and the system enters split brain or stops operating.