-
Notifications
You must be signed in to change notification settings - Fork 0
Concurrency
Vyra’s concurrency architecture is designed around one primary safety rule: user code must never execute on transport network I/O threads.
To ensure high network throughput and prevent deadlocks or thread starvation, Vyra cleanly isolates threads into three distinct layers:
flowchart TD
subgraph Layer1 ["Layer 1: transport I/O threads"]
AsyncIO["Async I/O event loops<br/>(non-blocking socket reads & writes)"]
Blocking["Dedicated blocking threads<br/>(one per worker queue)"]
end
subgraph Layer2 ["Layer 2: core processing"]
Dispatcher["MessageDispatcher<br/>(lightweight routing & classification)"]
end
subgraph Layer3 ["Layer 3: handler executor"]
SerDeser["Envelope & payload deserialization"]
UserHandler["User request handlers<br/>(handle)"]
UserSubscriber["User event consumers<br/>(subscribe)"]
CallbackCompletion["CompletionStage callbacks<br/>(thenAccept, whenComplete)"]
end
AsyncIO -->|raw byte envelope| Dispatcher
Blocking -->|raw byte envelope| Dispatcher
Dispatcher -->|submit task| SerDeser
SerDeser --> UserHandler
SerDeser --> UserSubscriber
SerDeser --> CallbackCompletion
classDef layer1 fill:#e3f2fd,stroke:#1565c0,color:#0d47a1
classDef layer2 fill:#fff3e0,stroke:#e65100,color:#bf360c
classDef layer3 fill:#e8f5e9,stroke:#2e7d32,color:#1b5e20
class Layer1 layer1
class Layer2 layer2
class Layer3 layer3
Three layers, one rule: transport threads only move bytes, user code always runs on the handler executor.
- Async I/O event loops: handle non-blocking socket reads and writes.
-
Dedicated blocking threads: each subscribed
REQUESTworker queue runs a single dedicated loop blocking on dequeue. - Guarantee: these threads only read raw bytes from the network and immediately pass them to Core. They never execute user logic, reflection, or deserialization.
- Deserialization of the message envelope is submitted immediately to the handler executor. Because your serializer is user code, running it on an I/O thread could introduce blocking or memory overhead.
-
MessageDispatcherinspects the envelope's kind and target to route it safely.
- User request handlers (
handle(...)), event subscribers (subscribe(...)), and callbacks attached to returnedCompletionStages execute here. - Even if your handler performs a blocking database query, takes seconds to compute, or crashes, the network I/O loops remain responsive.
What happens when incoming requests overwhelm your handler executor?
Instead of crashing the transport network thread or piling up unbounded memory:
- The executor rejects the submitted task.
- For requests: core immediately crafts an error response envelope with
RemoteError.Code.SERVER_BUSYand sends it back to the client. The client receives aVyraRemoteExceptionwithgetCode() == SERVER_BUSY. - For events: the dropped event is logged as dropped, and the transport thread continues uninterrupted.
This mechanism ensures graceful load shedding under traffic spikes.
By default, Vyra creates an internal daemon thread pool. However, for production workloads, you can provide your own ExecutorService:
// Example: custom thread pool optimized for blocking I/O handlers
ExecutorService handlerPool = new ThreadPoolExecutor(
16, 64,
60L, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(1000),
new NamedThreadFactory("vyra-handler-"),
new ThreadPoolExecutor.AbortPolicy()
);
Vyra vyra = RedisVyra.redis(client)
.serializer(JacksonSerializer.jackson())
.handlerExecutor(handlerPool)
.build();Tip
If you are on Java 21+, you can supply Executors.newVirtualThreadPerTaskExecutor() as your handler executor to handle massive concurrency for I/O-bound handlers without thread pool tuning.
-
Concurrent Vyra methods: all public methods on
Vyra(request,broadcastRequest,handle,publish,subscribe,register, andclose) are fully thread-safe and can be invoked concurrently from any thread. -
Envelope immutability: the
Messagerecord is strictly immutable. Headers are wrapped in unmodifiable maps. -
Serializer thread safety: all serializers provided by Vyra (
JacksonSerializer,GsonSerializer) are thread-safe. Custom serializers must also be thread-safe. - Pending requests: outbound requests are tracked in lock-free concurrent hash maps, ensuring fast correlation without synchronization bottlenecks.
- Learn how timeouts and thread-pool saturation (
SERVER_BUSY) are handled in Timeouts & errors. - Understand handler registration and options in Configuration.
Vyra documentation