Explain how Consul keeps its service catalog and KV store consistent, and how it tracks cluster membership — i.e. Raft vs gossip.
answer
- Raft = servers, consistent catalog/KV, quorum writes
- gossip = Serf/SWIM, all agents, membership+failure
- two pools: LAN + WAN
- clients forward RPC, hold no state
- 3 or 5 servers, odd for quorum
basics
~20 sConsul servers use the Raft consensus algorithm to keep one consistent, replicated copy of the catalog and KV store (writes go through an elected leader). Membership and failure detection use a gossip protocol (Serf/SWIM) across all agents.
solid answer
~50 sConsul splits two concerns. **Consistency of state** (the service catalog and KV store) is handled by **Raft** among the *server* agents: they elect a leader, and every write is replicated to a quorum of the server peers before being committed — giving a single, strongly consistent, linearizable copy of the data. A cluster of N servers tolerates ⌊(N-1)/2⌋ failures, which is why you run 3 or 5 servers. **Membership and failure detection** use a **gossip protocol** (HashiCorp's Serf, based on SWIM): every agent — clients and servers — periodically exchanges small UDP messages with random peers to learn who's alive, propagate join/leave events, and detect failures without a central coordinator. So gossip answers 'which nodes exist and are up?' quickly and scalably, while Raft answers 'what is the authoritative catalog/KV state?' consistently. Client agents forward reads/writes via RPC to the servers.
go deeper
Know at a high level: servers agree on data via Raft; all nodes track membership via gossip.
Distinguish which plane handles catalog/KV (Raft) vs membership/failure (gossip) and why quorum needs odd server counts.
Explain quorum fault tolerance, consistency modes, LAN/WAN pools, and client-agent RPC forwarding.
Reason about partition behavior, availability-vs-consistency boundaries, multi-DC federation, cluster sizing, and read-consistency tuning under load.
**Two independent planes.** Consul deliberately separates *consensus* (agreeing on authoritative data) from *membership* (knowing which nodes are in the cluster and alive). They use different algorithms with different consistency/latency trade-offs. **Raft (consensus, servers only).** Raft is a leader-based consensus algorithm. Among the **server agents**, one is elected **leader**; the rest are **followers**. All writes — service (de)registrations that reach the catalog, KV puts, ACL changes — are routed to the leader, appended to a replicated log, and are **committed only once a majority (quorum) of servers have persisted them**. This yields **linearizable, strongly consistent** state: a successful write is immediately visible to consistent reads. Fault tolerance is quorum-based: a 3-server cluster survives 1 failure, 5 survives 2. If quorum is lost, the cluster can't commit writes (it favors consistency over availability for the catalog/KV). Followers can serve **stale** reads for lower latency if you opt in (consistency modes: `default`, `consistent`, `stale`). **Gossip (membership + failure detection, all agents).** Membership uses **Serf**, HashiCorp's implementation of the **SWIM** protocol (Scalable Weakly-consistent Infection-style process group Membership). Every agent — both clients and servers — participates. Nodes periodically pick random peers and exchange compact messages over **UDP (with TCP fallback)** to: - **detect failures**: a node probes a random peer; if no ack, it asks other nodes to probe indirectly before declaring it `failed` (reduces false positives from a single bad link); - **disseminate events**: joins, leaves, failures, and user events spread epidemically ('infection-style'), reaching the whole cluster in O(log N) rounds; - **maintain an eventually-consistent** view of membership — fast and highly scalable, but *not* linearizable. There are actually **two gossip pools**: a **LAN pool** (all agents within a datacenter, low-latency failure detection) and a **WAN pool** (only server agents across datacenters, for multi-DC federation). **How they interact.** Client agents hold no catalog/KV state; they use gossip for membership and **RPC-forward** discovery queries and KV operations to the servers, which serve them from the Raft-backed store. When gossip detects a node failed, associated health checks go critical and the catalog reflects it. So: *gossip is the fast, scalable liveness/membership signal; Raft is the slower, consistent source of truth.* **Why the split matters for Spring apps.** - **Catalog reads are consistent** (Raft), so discovery won't hand you an instance the leader never registered; but you can trade consistency for latency with stale reads. - **Failure detection is quick and decentralized** (gossip), so a crashed node is noticed cluster-wide without hammering the leader. - **Availability boundary**: if the server quorum is lost, *writes* (new registrations, KV changes) stall even though gossip still 'sees' nodes — an operational gotcha. **Edge cases / gotchas.** - **Split-brain protection**: Raft requires a majority, so a network partition leaves at most one side able to commit — no divergent catalogs. - **Cluster sizing**: even numbers of servers add no fault tolerance and waste replication; use odd counts (3/5). - **UDP blocked**: gossip degrades if UDP is firewalled (falls back to TCP, but tune ports). - **Anti-entropy**: each agent periodically reconciles its local services/checks with the catalog, self-healing drift. **When this knowledge is used.** Sizing a Consul cluster, reasoning about behavior during partitions, choosing read consistency modes for discovery, and explaining why a registration succeeded but a query on another DC lagged (WAN gossip + non-replicated catalogs per DC).
- Why run 3 or 5 Consul servers rather than 2 or 4?Raft needs a majority quorum. Odd counts maximize fault tolerance per node: 3 tolerates 1 failure, 5 tolerates 2. An even count (e.g. 4) tolerates the same failures as 3 but needs an extra machine and raises the quorum, so it's wasteful.
- What consistency modes can a discovery/KV read use, and what's the trade-off?`consistent` forces a leader round-trip for linearizable reads (higher latency, no stale data); `stale` lets any server answer from its local log (lowest latency, may be slightly behind); `default` is leader-based but may serve slightly stale on leader lease. Choose stale for high-read discovery where tiny staleness is fine.
- How does multi-datacenter awareness work?Each datacenter runs its own Raft server cluster with its own catalog; servers additionally join a WAN gossip pool so datacenters know each other and can forward cross-DC queries. Data is not Raft-replicated across DCs — each DC is authoritative for itself.
saying these in an interview costs you the question
- Claiming client agents participate in Raft (only servers do)
- Saying gossip provides the strongly-consistent catalog (that's Raft)
- Believing KV is eventually consistent like DNS (it's Raft-linearizable)
- Thinking a 4-server cluster is more fault-tolerant than 3