DEV Community

Fred Feng
Fred Feng

Posted on

Openspreader: The missing piece of Java concurrency. `java.util.concurrent`, scoped to the cluster.

Thirteen coordination primitives as one Spring Boot starter. The APIs you already know,
spanning every instance of your application instead of one JVM. Zero lines of configuration to
start.

// A lock across every replica of this application
ProcessingMutex mutex = syncs.applicationMutex("daily-report");
if (mutex.tryAcquire()) {
    try { generateReport(); } finally { mutex.release(); }
}

// A scheduled task that runs on one instance per round. One annotation.
@Scheduled(cron = "0 0 2 * * *")
@MultiProcessingScheduled
public void syncOrdersDaily() { ... }
Enter fullscreen mode Exit fullscreen mode

What Problem Does It Solve?

synchronized only covers one JVM. Semaphore only limits one process. @Scheduled on three
replicas runs three times. The usual fix is Redis for a distributed lock, ZooKeeper for
coordination, and hand-written leader election for the scheduler: three pieces of
infrastructure, so that three replicas can agree on who does what.

openspreader does it with one dependency. Locks acquire in about a millisecond, cache reads are
served from local memory at over four million per second and never touch the network, and a
task dispatched to a peer runs there without your code knowing it moved.

Quick Start

<dependency>
    <groupId>com.chaconne-ai</groupId>
    <artifactId>openspreader</artifactId>
    <version>1.0.0-SNAPSHOT</version>
</dependency>
Enter fullscreen mode Exit fullscreen mode
<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

Not a line of configuration is required. Discovery defaults to 127.0.0.1, so two instances
on one machine with different server.port values form a cluster by themselves. Local
debugging works out of the box.

One line has to change for production:

spring.spreader.ip-addresses=192.168.0.111,192.168.0.63,192.168.0.77
spring.spreader.name=order-cluster
spring.spreader.multiprocessing.mutex.enabled=true
Enter fullscreen mode Exit fullscreen mode

With three instances started:

Node started: cluster=order-cluster, name=order-service, clusterPort=22000
Mutex service started

[instance 1] acquired nightly-report, generating
[instance 2] tryAcquire returned false, skipping this round
[instance 3] tryAcquire returned false, skipping this round
Enter fullscreen mode Exit fullscreen mode

Every component is off by default, so adding the starter to a running service changes
nothing until you ask for one.

Requirements

Java 17 or later
Spring Boot 4.1, built and tested against it
Optional Micrometer for Prometheus, Kryo for faster serialisation, spring-data-redis for cache overflow, Netty/MINA/Grizzly for an alternative transport
Ports one cluster port (22000 by default, identical on every node) plus one work port per node

How It Works

Everything rests on spreader: the cluster runs inside the same JVM as
your beans, so there is no broker and no registry.

  ┌─────────────────── your application ───────────────────┐
  │                                                        │
  │   @Scheduled   mutex.tryAcquire()   cache.get(k)       │
  │        │              │                  │             │
  │   ┌────┴──────────────┴──────────────────┴─────────┐   │
  │   │            openspreader components             │   │
  │   │  lock  semaphore  latch  cache  pool  rpc  dag │   │
  │   └────────────────────┬───────────────────────────┘   │
  │                        │                               │
  │   ┌────────────────────┴───────────────────────────┐   │
  │   │      spreader: gossip, members, leader         │   │
  │   └────────────────────┬───────────────────────────┘   │
  └────────────────────────┼───────────────────────────────┘
                           │ TCP/UDP
              ┌────────────┼────────────┐
         other replica           other replica
Enter fullscreen mode Exit fullscreen mode

Writes go through the leader, reads stay local. The registers for locks, permits, latches
and barriers live in the leader's memory; the cache replicates operations rather than data.

cache.incr("views", 1) on a follower
    │
    ├─ forwarded to the leader
    │
    ├─ leader applies it, assigns version 1007
    │
    └─ broadcasts {v1007, INCR, "views", 1} to every replica
                      │
         each replica applies it to its OWN copy
Enter fullscreen mode Exit fullscreen mode

That is why setbit on a 100MB bitmap ships one datagram rather than the bitmap, and why a
cluster-wide Bloom filter is practical at all.

Code Examples

Five representative components. Every component has a runnable example under
com.chaconneai.openspreader.example.

Lock

Input

ProcessingMutex mutex = syncs.applicationMutex("bulk-export");
Enter fullscreen mode Exit fullscreen mode

Execution: three ways to take it, matching three attitudes.

mutex.tryAcquire();                    // skip if taken, do not wait
mutex.acquire(3, TimeUnit.SECONDS);    // wait three seconds, then skip
mutex.acquire();                       // wait indefinitely
Enter fullscreen mode Exit fullscreen mode

Output

mutex.isHeld();              // is anyone holding it
mutex.currentOccupied();     // WHO holds it, which answers "why did my task not run"
mutex.release(5_000);        // release, but keep others out for five seconds
Enter fullscreen mode Exit fullscreen mode

That cooldown solves a real problem: a task finishes, another replica immediately takes the
lock and runs it again. A cooldown window removes the duplicate.

Semaphore

Input: protect a third party API at five concurrent calls.

ProcessingSemaphore sem = syncs.applicationSemaphore("third-party-api", 5);
Enter fullscreen mode Exit fullscreen mode

Execution

if (sem.tryAcquire(3, TimeUnit.SECONDS)) {
    try { callThirdParty(); } finally { sem.release(); }
}
Enter fullscreen mode Exit fullscreen mode

Output: the five permits are shared cluster-wide, whatever the replica count. That is
exactly what a single process Semaphore cannot do: three replicas each holding a local
semaphore of 5 gives you an actual concurrency of 15.

Cache

Input: 50 commands over four data structures, plus bitmaps and Bloom filters.

cache.set("k", bytes, 10, TimeUnit.MINUTES);
cache.hset("user:1", "email", bytes);
cache.zadd("leaderboard", member, 99.5);
cache.setbit("seen", offset, true);
cache.zrangeByScore("series", from, to, 0, 500);     // a page, not a million members
cache.zremrangeByScore("series", 0, cutoff);         // a retention policy in one operation
Enter fullscreen mode Exit fullscreen mode

Execution: reads come from local memory; writes are forwarded to the leader and broadcast
back incrementally.

Output

reads     over 4,000,000 per second, never touching the network
writes    roughly 2,000 per second (TCP), 8,300 (UDP)
Enter fullscreen mode Exit fullscreen mode

So it suits shared state that is read far more than written: configuration, allow-lists,
rate-limit counters, sessions. A write becomes visible elsewhere after milliseconds. It is
eventually consistent, not a database.

When memory is not enough: register a CacheStore bean and eviction becomes a move
rather than a loss.

spring.spreader.multiprocessing.cache.external.enabled=true
Enter fullscreen mode Exit fullscreen mode
A key lives in exactly one half Memory, or outside. Reads ask memory first; writes go where the key already is
Eviction moves instead of deleting Same sampling, same policy, but the key is written across before the memory copy goes
Only the leader writes outside A follower replaying the replication stream touches memory only
Bitmaps stay in memory A Bloom filter lookup does seven bit tests, which outside would be seven round trips
Nothing is enabled by default Without the bean, store is the local store: same object, no wrapper, no extra call

MapReduce

Input: implement a job, register it as a bean.

@Component("wordCount")
public class WordCountJob implements MapReduceJob<String, String, Integer, Integer> {

    public List<String> split(String text, int suggestedShards) {
        return splitByLines(text, suggestedShards);
    }

    public void map(String shard, Emitter<String, Integer> emitter) {
        for (String word : shard.split("\\W+")) {
            if (!word.isBlank()) {
                emitter.emit(word.toLowerCase(), 1);
            }
        }
    }

    public Integer reduce(String word, List<Integer> counts) {
        return counts.stream().mapToInt(Integer::intValue).sum();
    }
}
Enter fullscreen mode Exit fullscreen mode

Execution: where each phase runs.

split     the submitter only, once
map       every node, the submitter included, once per shard
reduce    wherever the key's hash sends it, once per distinct key
Enter fullscreen mode Exit fullscreen mode

Output

MapReduceResult<String, Integer> result = mapReduce
        .<String, String, Integer, Integer>submit("wordCount", input)
        .get(5, TimeUnit.MINUTES);
Enter fullscreen mode Exit fullscreen mode

Combining the partial results is not yours to write: one key is reduced on exactly one node.

Every node needs this bean, at the same version. Half the nodes on a new map and half
on the old one produce a wrong answer rather than an error. Do not submit jobs during a
rolling deployment.

RPC

Input: enable the scan, declare an interface pointing at another application's name.

@SpringBootApplication
@EnableRpcClients(basePackages = "com.example.client")
public class Application { }

@RpcClient(serviceId = "inventory-service",
           timeout = 3000,
           maxConcurrent = 50,                       // in-flight ceiling, fail fast when full
           fallback = InventoryFallback.class)
public interface InventoryApi {

    int stockOf(String sku);
}
Enter fullscreen mode Exit fullscreen mode

Execution: inject and call it like any bean. The server side needs a bean with a matching
method carrying @MultiProcessingCall.

int stock = inventory.stockOf("widget-1");
Enter fullscreen mode Exit fullscreen mode

Output: routing is by application name rather than URL, and the target is chosen afresh on
every call, so scaling on the other side takes effect immediately.

Setting Why it matters
maxConcurrent Without it, requests pile up locally when the downstream slows, and what falls over is your own process
maxRetries A failure is retried on a different instance each time, so the called method should be idempotent. Retrying after a timeout is ambiguous: the request may have completed with only the reply lost

Configuration

Property Default Description
spring.spreader.name default Cluster name, the only means of isolation between dev, staging and prod
spring.spreader.ip-addresses 127.0.0.1 Where to look for peers
spring.spreader.advertise-host auto Set this in containers
...cache.max-keys 1000000 Per process, not per cluster: every node holds a full replica
...cache.eviction-policy LRU Samples 5 keys and evicts the least recently used, as Redis does. Exact LRU would need a global lock on every read
...cache.external.enabled false Build the Redis-backed overflow store
...mutex.request-timeout-ms 3000 Shorter than the cache's: a lock request not answered quickly is better retried
...mutex.lease-ms 15000 Renewed while held, so a task may run as long as it likes. Only a crashed holder loses it

Performance

Conditions: single machine over loopback, 3-node cluster in one JVM, 4-core container, JDK
serialization. Absolute values do not transfer; the ratios do.

Reads never leave the process

Operation QPS
exists 6,562,196
size 4,566,915
stats (aggregate, four numbers) 4,130,498
hget (single field) 2,104,340
getbit 1,548,987
zscore 1,210,631
lrange (50 elements) 345,364
zrange (50 members) 287,519
hgetAll (all fields) 68,489

hget against hgetAll is 30x. Calling hgetAll in a loop is the easiest accidental
mistake in this API, and it costs more than the read/write gap does.

Writes go through the leader

Operation QPS
max (aggregate) 7,314
incr 6,524
min / sum (aggregate) 6,356 to 6,539
zadd 5,504
rpush 4,500
hset 4,080
setbit (fixed capacity) 3,419
setbit (growing) 1,991

Aggregates are the fastest writes. max/min/sum keep four numbers, so there is no
allocation, no resizing and no hash lookup, and the value travels in the arg field rather
than as a payload. That is what makes them suitable for downsampling a time series.

The ceiling, and where it comes from

Writes are serialised globally: every write goes through the leader and takes one state lock,
because a monotonic version number is what lets all replicas replay in the same order. The
lock's granularity is the version number's granularity.

Read scaling Linear with node count. Reads are local
Write scaling Flat. Adding nodes does not raise write throughput; it raises the leader's broadcast fan-out
Practical write ceiling ~2,000/s on TCP, ~8,300/s on UDP

So the question to ask is how far your write rate is from those numbers. Far means no problem.
Close means that data may not belong here.

Limitations & Trade-offs

Every coordination capability rests on one thing: whoever holds the cluster port is the
leader
, and the write path for locks, permits, latches, barriers and the cache is all there.

Premise Consequence
Leader uniqueness rests on timing, not consensus Under partition each side may have a leader, and two parties hold the same lock. Fine for work where a repeat is wasteful: scheduled tasks, cache warming, batch de-duplication. Not for a transfer or a debit, which belongs in a database transaction or behind an idempotence key
A change of leader leaves a gap Takeover measured at 3.3 to 4.4 seconds, during which lock acquisitions and cache writes retry. Leave business timeouts headroom, or you will see "it failed once and was fine a few seconds later"
Limit Detail
The cache is a cache Every node holds a full replica in heap. Not a database, not durable
Eviction is local Nodes can disagree about which keys are resident
Everything is in-process No persistence by default: a full cluster restart starts from empty
Resume is at-least-once A DAG node that ran and reported back, whose record was not yet written when the coordinator died, is a step that happened and is not in the record
Cross-machine performance is unmeasured Every number above is one JVM over loopback

Summary

  1. Thirteen coordination primitives, one starter, zero configuration to start
  2. The APIs you already know, scoped to every instance of the application rather than one JVM
  3. Reads never leave the process: over four million per second, because every node holds a full cache replica
  4. Writes have a clear ceiling: ~2,000/s on TCP, from a monotonic version number that must be assigned serially
  5. Aggregates are the fastest writes (7,314/s), four numbers in constant memory, suited to downsampling
  6. hget is 30x faster than hgetAll, a trap easier to fall into than the read/write gap
  7. Every component is off by default, so adding the starter changes no behaviour
  8. Memory can overflow to an external store, turning eviction from a loss into a move. Off by default
  9. Two premises to accept: leader uniqueness rests on timing not consensus, and a change of leader leaves a gap
  10. The cluster layer underneath is spreader, where the same post builds a replicated H2 on it

Top comments (0)