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