Configure a Vert.x 4.x Clustered Event Bus with Custom Codecs and Reliable Request-Response
Step‑by‑step guide to start a Vert.x 4.x clustered event bus, register a custom codec, and use timed request‑response messaging with failure handling.
16 Jul 2026, 04:43 UTC

Overview
This guide shows how to run a Vert.x application in clustered mode, register a custom codec once per node, and use request‑response messaging with explicit timeouts. The desired outcome is two (or more) Vert.x instances that can exchange typed objects reliably, detect failures, and recover without manual intervention.
Prerequisites
- Java 11+ installed and
JAVA_HOME correctly. - Vert.x 4.x CLI (
vertx) available inPATH. - Network access between the hosts that will run the nodes (default Hazelcast ports 5701‑5708 open).
- A simple POJO you want to send, e.g.,
com.example.Messagewith fieldsString idandint value.
Step 1 – Prepare a Hazelcast configuration
Create a file hazelcast.xml that disables multicast (unsuitable for most clouds) and enables TCP/IP discovery.
<hazelcast xsi:schemaLocation="http://www.hazelcast.com/schema/config hazelcast-config-5.0.xsd"
xmlns="http://www.hazelcast.com/schema/config"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance">
<network>
<join>
<multicast enabled="false" />
<tcpip enabled="true">
<member>10.0.0.11</member>
<member>10.0.0.12</member>
</tcpip>
</join>
</network>
</hazelcast>
Place this file in the working directory of each node or reference it with an absolute path.
Step 2 – Register the custom codec identically on every node
Codec registration must happen before any message is sent and must be identical across the cluster; otherwise a ClassNotFoundException occurs on deserialization.
public class MessageCodec implements io.vertx.core.buffer.MessageCodec<Message, Message> {
@Override
public void encodeToBuffer(io.vertx.core.buffer.Buffer buf, Message msg) {
buf.appendString(msg.getId()).appendInt(msg.getValue());
}
@Override
public Message decodeFromBuffer(int pos, io.vertx.core.buffer.Buffer buf) {
String id = buf.getString(pos);
int pos2 = pos + id.getBytes(StandardCharsets.UTF_8).length + 2; // length prefix omitted for brevity
int value = buf.getInt(pos2);
return new Message(id, value);
}
@Override
public Message transform(Message msg) { return msg; }
@Override
public String name() { return "message-codec"; }
@Override
public byte systemCodecID() { return -1; }
}
// Registration (call once per Vertx instance, e.g., in a startup verticle)
vertx.eventBus().registerDefaultCodec(Message.class, new MessageCodec());
Deploy the same verticle containing this registration on each node before any business verticles start.
Step 3 – Launch clustered Vert.x instances
Run the application with the -cluster flag and point Vert.x to the Hazelcast config.
# Run on node 1
vertx run com.example.MainVerticle -cluster -conf hazelcast.xml
# Run on node 2 (same command, different host)
vertx run com.example.MainVerticle -cluster -conf hazelcast.xml
Required permissions: the user must be able to open TCP ports 5701‑5708 and create temporary files in the working directory.
Step 4 – Use request‑response with explicit timeouts
When a verticle needs a reply, wrap the call in DeliveryOptions that set a send timeout and specify the codec name (optional if the default codec is registered).
DeliveryOptions opts = new DeliveryOptions()
.setSendTimeout(5000) // wait max 5 s for a reply
.setCodecName("message-codec");
EventBus eb = vertx.eventBus();
eb.request("service.address", new Message("req-1", 42), opts, reply -> {
if (reply.succeeded()) {
Message resp = reply.result().body();
System.out.println("Got " + resp.getValue());
} else {
System.err.println("Request failed: " + reply.cause());
}
});
The temporary reply handler created by request() is automatically cleaned up when the timeout expires or a reply is received, preventing leaks.
Step 5 – Verify cluster formation and messaging
After starting both nodes, look for log lines similar to:
INFO io.vertx.core.impl.clusteredimpl.HazelcastClusterManager - Cluster createdINFO io.vertx.core.impl.VertxImpl - Node address: /10.0.0.11:5701INFO io.vertx.core.impl.VertxImpl - Node UUID: a1b2c3d4‑e5f6‑7890‑abcd‑ef1234567890
To test messaging, send a request from one node and assert the reply on the other (example pseudo‑code, not actual output):
// On node A
vertx.eventBus().request("service.address", new Message("ping", 7), new DeliveryOptions().setSendTimeout(3000), ar -> {
if (ar.succeeded()) {
System.out.println("Reply: " + ar.result().body().getValue());
} else {
ar.cause().printStackTrace();
}
});
Expected check: the receiving node logs the request, processes it, and sends a reply; the sending node logs the reply within the timeout window.
Step 6 – Failure detection and recovery
Monitor cluster health and message failures:
- Add an exception handler to consumers to catch serialization or processing errors.
- Use the completion handler of
Vertx.clusteredVertx()to react to cluster loss. - Enable Hazelcast’s split‑brain detection (via
split-brain-protectioninhazelcast.xml) if your deployment tolerates brief partitions.
vertx.eventBus().consumer("service.address").exceptionHandler(err -> {
System.err.println("Consumer error: " + err);
// optionally trigger alert or fallback logic
});
Vertx.clusteredVertx(new VertxOptions().setHAEnabled(true), res -> {
if (res.succeeded()) {
Vertx clustered = res.result();
clustered.close().onComplete(v -> System.out.println("Cluster node shut down cleanly"));
} else {
System.err.println("Failed to form cluster: " + res.cause());
}
});
Rollback procedure
If you need to revert to a non‑clustered setup:
- Stop all running Vert.x instances (
Ctrl+Cin each terminal orkillthe Java process). - Remove the
-clusterflag from your launch command. - Optionally rename or delete
hazelcast.xmlso Vert.x falls back to the default local-only event bus. - Restart the verticles; they will now operate in a single‑JVM mode.
Limitations and practical verification
- Ordering is guaranteed only per‑sender‑per‑address; interleaving across different addresses may occur.
- The event bus does not persist messages; for durability integrate an external broker (Kafka, RabbitMQ) via Vert.x clients.
- Custom codecs must be stateless and thread‑safe because the same instance can be invoked from multiple event‑loop threads.
- To verify codec correctness, send a known object and assert that the received object’s fields match the original values.
- To test network‑partition resilience, block Hazelcast ports between nodes with a firewall rule, confirm that
eventBus.clustered()returnsfalseand that local‑only messages (DeliveryOptions.setLocalOnly(true)) still succeed.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.