Apache Iggy
Clustering

Viewstamped Replication

How Apache Iggy replicates data across a cluster with VSR

Apache Iggy replicates clusters with Viewstamped Replication Revisited (VSR). VSR keeps an ordered state machine consistent across replicas and elects a new primary when the current primary fails.

Clustering is built into the standard iggy-server binary. There's no separate build, binary, or Cargo feature: the same server runs single-node or clustered, and cluster.enabled in the configuration decides which one you get. Clustering is disabled by default.

Read the VSR paper used by the project for the protocol specification.

Iggy hasn't reached a 1.0 production release yet. The cluster protocol and configuration described here can still change between pre-release versions.

Current implementation

The implementation includes:

  • normal operation with quorum commits
  • view changes and deterministic primary selection
  • metadata and partition replication
  • persistent WAL recovery and snapshots
  • replica rejoin with journal repair
  • client sessions, request fencing, and duplicate detection
  • replica authentication and optional TLS
  • leader redirection in the client SDKs
  • deterministic simulation and cross-SDK BDD coverage

The main components are:

ComponentSourceResponsibility
Servercore/serverSharded runtime, listeners, recovery, and request dispatch
Consensuscore/consensusVSR state machine, quorum, view changes, and timeouts
Wire protocolcore/binary_protocol/src/consensusFixed-size VSR headers and replica messages
Simulatorcore/simulatorDeterministic failures, delays, and network partitions
Rust SDKcore/sdkVSR framing, sessions, retries, and leader redirection

Replication model

Iggy splits replication by namespace:

PlaneConsensus groupReplicated work
MetadataOne group on shard 0Streams, topics, users, permissions, consumer groups, and access tokens
PartitionOne group per partitionMessages and consumer offsets

Each group can have a different primary. Metadata writes route through the metadata primary. Partition writes route through the primary for that partition.

Reads use the local replicated state where possible. Writes return after the required VSR commit or return a retryable error when the replica is changing view or catching up. See Client failover for how the SDKs handle both.

Failure handling

VSR uses three main flows:

  1. Normal operation: the primary assigns an operation number, sends Prepare, and commits after a quorum sends PrepareOk.
  2. View change: replicas exchange StartViewChange and DoViewChange, then the new primary sends StartView.
  3. Recovery: a restarted or lagging replica requests the current view and repairs missing WAL ranges before serving current state.

A cluster of 2f + 1 replicas tolerates f unavailable replicas. Use at least three replicas for one-node fault tolerance. A two-node cluster is useful for development, but it cannot make progress after either node fails.

Wire protocol and SDKs

VSR framing is the only wire protocol. Every SDK (Rust, Go, Java, C#, C++, Node.js, Python) speaks it, and every BDD suite runs against the clustered-capable iggy-server. Clients built before the VSR migration cannot talk to current servers.

The Rust SDK needs no feature flags. Depend on the published crate:

[dependencies]
iggy = "0.11.0-edge.4"
tokio = { version = "1", features = ["full"] }

or track the repository directly with iggy = { git = "https://github.com/apache/iggy" }.

The client API is the same in single-node and cluster mode:

use iggy::prelude::*;

#[tokio::main]
async fn main() -> Result<(), IggyError> {
    let client = IggyClientBuilder::from_connection_string(
        "iggy://iggy:iggy@127.0.0.1:8090",
    )?
    .build()?;

    client.connect().await?;

    let cluster = client.get_cluster_metadata().await?;
    println!("cluster: {}", cluster.name);

    Ok(())
}

Status

VSR started as a parallel server-ng implementation behind a vsr build feature. That phase is over. The VSR server moved into the mainline core/server crate and the iggy-server binary, the feature flag is gone, and the legacy server and its wire format were deleted. What this page describes is the only server implementation.

Next steps

  • Deploy a cluster: build the binaries and run a three-node cluster
  • Configuration: every [cluster] setting and its validation rules
  • Security: replica authentication, key rotation, and TLS
  • Client failover: reads, writes, and leader redirection from the client side

On this page