DEV Community

Fred Feng
Fred Feng

Posted on

Spreader: Your JVM processes are the cluster. Nothing else to install.

Turn plain JVM instances into a cluster that elects a leader, tracks every member and passes
messages between them. One dependency, one line of configuration, and no ZooKeeper, etcd,
Consul or registry anywhere in the picture. The processes are the cluster.

GossipCluster cluster = GossipCluster.create(
        GossipConfig.builder().ipAddresses("10.0.0.1", "10.0.0.2").build());
cluster.start();

if (cluster.isLeader()) {
    runTheNightlyJob();          // exactly one instance does this
}
Enter fullscreen mode Exit fullscreen mode

What Problem Does It Solve?

Three machines running the same application, and you need them to agree on which one runs the
nightly job, to notice within seconds when one dies, and to send each other a message. The
standard answer is a new piece of infrastructure: something to deploy, monitor, upgrade and get
paged about at 3am, for a problem that is not your product.

spreader keeps it inside the processes. Every node runs the same code and holds the same
complete member list, with no special node and no role to assign. Measured: 18,464 cross-node
round trips per second, and a dead leader replaced in 3 to 4 seconds.

Quick Start

<dependency>
    <groupId>com.chaconne-ai</groupId>
    <artifactId>spreader</artifactId>
    <version>1.0.0-SNAPSHOT</version>
</dependency>
Enter fullscreen mode Exit fullscreen mode

Snapshots need the repository declared too:

<repositories>
    <repository>
        <id>central-snapshots</id>
        <url>https://central.sonatype.com/repository/maven-snapshots/</url>
        <snapshots><enabled>true</enabled></snapshots>
    </repository>
</repositories>
Enter fullscreen mode Exit fullscreen mode

Open two or three terminals and run one in each:

java -cp spreader-1.0.0-SNAPSHOT.jar \
     com.chaconneai.spreader.example.BestPractice \
     --cluster=demo --port=22000 --ipAddresses=127.0.0.1 --portRetry=10
Enter fullscreen mode Exit fullscreen mode

What you see:

Node started: cluster=demo, id=a1b2c3d4, clusterPort=22000, transport=TCP/NioTcpTransport
joined: the only node so far, so this one leads

===== cluster [demo] =====
this node: demo-app@127.0.0.1:30001  [leader]
leader:    demo-app@127.0.0.1:30001
members:   1

>>> [joined] demo-app@127.0.0.1:30002, members now 2
>>> [joined] demo-app@127.0.0.1:30003, members now 3
Enter fullscreen mode Exit fullscreen mode

Kill the leader, and within seconds:

>>> [leader changed] 127.0.0.1:30001 -> 127.0.0.1:30002  (this node is now the leader)
Enter fullscreen mode Exit fullscreen mode

Requirements

Java 17 or later
Optional Netty, MINA or Grizzly, to replace the built-in NIO transport
Ports one cluster port (22000 by default, identical on every node) plus one work port per node, chosen automatically

How It Works

Two mechanisms, both deliberately boring.

Membership spreads by SWIM gossip. Every second each node picks a few peers at random and
exchanges member lists directly, never through a coordinator.

    Node A ◄──────── gossip ────────► Node B
       ▲  holds :22000                   ▲
       │  = is the leader                │
       └──────────── gossip ─────────────┘
                      │
                   Node C
Enter fullscreen mode Exit fullscreen mode

Leadership is pre-emptive. The cluster has one fixed port and whichever node holds it
leads. A new node scans the configured addresses, knocks on that port, and either finds a
leader or becomes one.

leader dies
    │
    ├─ failure detection notices (probe, indirect probe, suspect timeout)
    │
    ├─ every eligible node computes its rank in the takeover order
    │
    ├─ rank 0 waits 0ms, rank 1 waits takeoverDelayMs, rank 2 waits 2x ...
    │
    └─ first to bind :22000 is the new leader; the rest see it and stand down
Enter fullscreen mode Exit fullscreen mode

Pre-emptive and consensus-based election are two routes, not two grades

Where a leader's legitimacy comes from is the root of every other difference. Pre-emptive: from
holding a physically exclusive resource. Consensus-based: from a majority's authorisation,
since two majorities must intersect and a node in that intersection will not vote twice in one
term.

Pre-emptive (this library) Consensus-based (Raft)
How a leader arises whoever takes the port whoever a majority votes for
Must the leader be up to date? No. A node that just restarted can lead Yes. A stale log cannot win a vote
Terms none, so no built-in fencing monotonic; an older term is refused
One node works works
Two nodes, one dies the survivor takes over stalls, no majority
Membership changes soft state, gossiped, free configuration state, changed through the log
Disk none fsync before voting
Under partition both sides can lead and accept writes the minority elects nobody and refuses writes
Takeover time 3 to 4 seconds, dominated by failure detection typically a few hundred milliseconds

Two consequences shape how to use it:

  • A pre-emptive leader carries no data authority. It is not required to be up to date, so any state kept on it has to be rebuildable.
  • With no terms there is no built-in fencing. A leader already replaced can still write to an outside resource. Where that matters, have the leader issue a monotonic token, carry it into whatever you write, and have that resource refuse older ones.

The dividing line is the nature of the work, not its importance. Coordination where a
repeat is merely wasteful fits this route. Work where a repeat causes real harm needs
consensus-based election.

Code Examples

Example 1: members and the leader

Input

List<Node> all = cluster.members();                        // the whole cluster
List<Node> peers = cluster.membersOf("order-service");     // just my own application
Node leader = cluster.leader();                            // can be null
Enter fullscreen mode Exit fullscreen mode

Execution: reach for membersOf more often than members. "The other replicas of my own
application" is usually the question actually being asked, especially once a cluster runs more
than one application.

Output: each Node carries what you need.

node.address();     // an address that can be dialled
node.name();        // application name
node.state();       // ALIVE / SUSPECT / LEFT
node.metadata();    // your own business metadata, gossiped cluster-wide
Enter fullscreen mode Exit fullscreen mode

metadata is an extension channel: to let everyone know this node's version, zone or canary
flag, put it in the configuration with .metadata("zone", "az-1") and every node reads it off
the member list.

Example 2: a real trap

Input: three replicas, and a scheduled task that should run once.

// Wrong
if (cluster.isLeader()) {
    syncOrdersDaily();
}
Enter fullscreen mode Exit fullscreen mode

Execution: the leader is a cluster-wide notion that does not distinguish applications.
If your cluster also runs an api-facade and the cluster port happens to be held by one of its
instances...

Output: every replica of order-service sees isLeader() == false and the task never runs
at all
.

The right tool is openspreader's @MultiProcessingScheduled, which picks one executor per
application. Or elect within your own application against membersOf(self().name()).

Example 3: a replicated H2

Input: any node writes, and callers do not need to know whether they are the leader.

store.write("INSERT INTO product (name, price) VALUES ('widget-1', 10.99)");
Enter fullscreen mode Exit fullscreen mode

Execution: the write takes one path only.

write on a follower          write on the leader
       │                            │
 sendToLeader(sql)                  │
       │                            │
       └────► the leader applies it locally
                      │
             multicast to every replica
                      │
             each applies it to its own H2
Enter fullscreen mode Exit fullscreen mode
public void write(String sql) throws SQLException {
    byte[] message = encode(cluster.isLeader() ? REPLICATE : WRITE_REQUEST, sql);
    if (cluster.isLeader()) {
        // Apply first, then tell the others. The other order would have a replica
        // holding a row the leader itself does not
        apply(sql);
        cluster.multicastOn(CHANNEL, cluster.self().name(), message, false);
    } else if (!cluster.sendToLeaderOn(CHANNEL, message)) {
        // No leader at this instant. Surfacing it beats silently dropping the write
        throw new SQLException("No leader to take the write; retry in a moment");
    }
}
Enter fullscreen mode Exit fullscreen mode

Output: every node reads its own copy, and no read leaves the process.

java -cp spreader.jar:h2.jar \
     com.chaconneai.spreader.example.ReplicatedH2Example --port=22000
Enter fullscreen mode Exit fullscreen mode
[127.0.0.1:30001 leader] rows=8 applied=8 members=3
[127.0.0.1:30002]        rows=8 applied=8 members=3
[127.0.0.1:30003]        rows=8 applied=8 members=3
Enter fullscreen mode Exit fullscreen mode

One step is easy to miss: a node joining an established cluster starts empty. Replicating
only the increments leaves it permanently short of history, answering reads with nothing, which
is worse than answering slowly. So ask for a full copy on join:

@Override
public void onClusterJoined(Node self, boolean alone) {
    if (!alone) {
        cluster.sendToLeaderOn(CHANNEL, encode(SNAPSHOT_REQUEST, ""));
    }
}
Enter fullscreen mode Exit fullscreen mode

H2 has a convenient statement for this: SCRIPT dumps the whole database as SQL, one statement
a row. The leader dumps it and sends it over, the new node runs DROP ALL OBJECTS and replays
it, and it is caught up.

The complete code is in com.chaconneai.spreader.example.ReplicatedH2Example, about 350 lines.
It needs no H2 at compile time, using only java.sql from the JDK, so the spreader jar
does not carry it.

Configuration

Defaults are tuned for a small cluster on a reliable LAN. One line genuinely has to change.

Property Default Description
ip-addresses 127.0.0.1 Where to look for peers. A few seeds are enough; a new node learns the full member list from any one of them
cluster-name default The only means of isolation. Nodes with different names ignore each other even when the network can reach them
node-name default The application name. Instances sharing one are replicas, and the layers above use it to tell who belongs with whom
cluster-port 22000 The port that decides leadership. Identical on every node
advertise-host auto-detected Set this in containers. Auto-detection may pick an address no peer can reach
leader-eligible true false makes this application join and work but never contend for leadership
transport-type TCP TCP unless you have measured otherwise
payload-ack true Wait for an ACK and resend. Off is roughly 3x the throughput, since the ACK round trip dominates small-message cost
probe-timeout-ms 800 Raise this first on a congested network. Too low and healthy nodes get suspected during a GC pause

Performance

Numbers expose a property of the design. Absolute values will not transfer to your
hardware
: every cross-node call here costs a loopback hop instead of a real network one. The
ratios will.

Conditions: single machine over loopback, 8 threads x 2000 synchronous round trips, ~300
byte payloads, JDK serialization, 4-core container.

Transport x protocol: the pairing is the unit of choice

Combination QPS avg P50 P99 max
NIO / TCP 18,464 0.426 ms 0.367 ms 1.69 ms 5.4 ms
NIO / UDP 13,023 0.606 ms 0.423 ms 6.46 ms 15.0 ms
MINA / TCP 12,249 0.640 ms 0.438 ms 5.18 ms 15.4 ms
NETTY / UDP 10,682 0.746 ms 0.486 ms 7.19 ms 28.9 ms
MINA / UDP 10,494 0.751 ms 0.464 ms 7.24 ms 21.8 ms
GRIZZLY / TCP 9,171 0.864 ms 0.580 ms 7.68 ms 21.3 ms
GRIZZLY / UDP 8,751 0.910 ms 0.644 ms 6.96 ms 20.2 ms
NETTY / TCP 7,081 1.115 ms 0.495 ms 14.46 ms 56.8 ms

Pick the cell, not the framework. Netty is second-best on UDP and last on TCP, with a P99
twice anyone else's. MINA is the opposite. "Which transport is fastest" has no answer.

Treat anything under ~20% as noise: an earlier run on a 12-core host put Netty/UDP on top, and
NIO/TCP measured 16,351 and 14,637 on two consecutive runs.

The same four transports, through the components

Operation (TCP) NIO NETTY MINA GRIZZLY
lock acquire + release 875 686 616 648
cache write (incr) 2,110 1,733 1,440 1,748
task dispatch to a peer 11,581 10,028 11,473 9,148
Operation (UDP) NIO NETTY MINA GRIZZLY
lock acquire + release 2,878 2,162 3,009 720
cache write (incr) 8,326 8,237 2,645 1,618
task dispatch to a peer 5,809 10,787 10,407 8,856

Two findings:

  • Coordination runs 3 to 4 times faster over UDP (locks 875 to 2,878, cache writes 2,110 to 8,326). These are short one-way messages plus a reply, and TCP's acknowledgement and congestion control are pure overhead on that shape. A sustained request/response stream is the opposite, and TCP wins there.
  • Grizzly collapses on UDP: 720 against NIO's 2,878 is four-fold, far outside the noise band, and it appears on every coordination row while leaving task dispatch alone.

Reads are absent from both tables on purpose: they come from local memory and never touch the
transport, so comparing them would only mislead.

Dispatch threads are not a throughput dial

Mode msg/s per message
4 producers, 4 dispatch threads 5,334,329 187 ns
serial dispatch (payload-dispatch-threads=1) 5,100,959 196 ns
concurrent dispatch (payload-dispatch-threads=4) 3,453,596 290 ns

Rows two and three share a single producer, and the serial dispatcher wins by almost 50%.
Coordination between dispatch threads costs more than the parallelism returns. Raise this
setting when several threads publish concurrently, leave it at 1 otherwise.

Limitations & Trade-offs

Limit Detail
Election is pre-emptive, not consensus-based Under partition two sides can each elect a leader. It heals when the partition does, and every occurrence is counted. Work where a repeat causes real harm does not belong on this route
Membership is eventually consistent Right after a change, different nodes briefly hold different views
UDP messages can be lost payload-ack covers it with resend and dedup, at a throughput cost
Cross-machine performance is unmeasured Every number above is one JVM over loopback

Where it does not fit: thousands of nodes, clusters spanning datacentres, or a leader that must
stay correct under partition. The first two make gossip convergence unpredictable; the third
calls for Raft.

Summary

  1. One jar, one line of configuration. No registry, no coordinator. The processes are the cluster
  2. Fully decentralized: every node runs the same code and holds the same complete member list
  3. Election needs no consensus round: whoever holds the cluster port leads, and the OS already guarantees one port, one process
  4. Pre-emptive and consensus-based are two routes, not two grades. The line falls on the nature of the work, not its importance
  5. 18,464 cross-node round trips per second measured, and a dead leader replaced in 3 to 4 seconds
  6. Pick the cell, not the framework: Netty is last on TCP and second on UDP. There is no "fastest transport"
  7. Coordination runs 3 to 4x faster over UDP, while a sustained request/response stream favours TCP
  8. The two failures that log nothing each have a counter: dropped messages and a second leader
  9. Limits are stated, not hidden: partition dual leaders, eventual consistency, UDP loss, cross-machine unmeasured
  10. For ready-made locks, caches and scheduled-task exclusion, see openspreader

Top comments (0)