-
Notifications
You must be signed in to change notification settings - Fork 0
Parallel Consumers
Introduced in 2.4.0. Consumers are sequential by default; executor = "" (the default) keeps the
existing inline path, so nothing changes for a bundle that does not opt in.
Public guide: docs.eventoframework.com → Parallel Consumers.
Register a bounded executor on the bundle builder, then name it on the handler:
.addConsumerExecutor(ConsumerExecutors.virtual("read-model", 64))@EventHandler(executor = "read-model", retry = 3)
void on(OrderTotalRecomputed event) { … }A name is a capacity budget shared bundle-wide — two handlers naming "read-model" draw from the
same 64 permits. Referencing an unregistered executor fails bundle start-up; it does not silently
fall back to sequential.
Factories: ConsumerExecutors.virtual(name, n), pooled(name, n), partitioned(name, lanes).
Not available on @SagaEventHandler.
The checkpoint advances when a task starts, not when it completes.
ConsumerExecutor.submit returns a future that completes at task start, and the consume loop awaits
that future before committing. This one choice does two jobs:
- It is the backpressure. The loop can never run more than the executor's capacity ahead of completion, so no unbounded prefix of the event store is ever enqueued.
- It bounds the crash-loss window to the set of running tasks, rather than to a queue depth.
Submission is time-boxed (5 s default). On saturation the cycle ends at the last started event and
releases the consumer lock — important, because the JDBC ConsumerLock pins a pooled connection
for the whole cycle, and parking on a saturated executor would hold that connection hostage.
| Guarantee | Under a consumer executor |
|---|---|
| Ordering | None between parallel events |
| Delivery |
At-most-once for events in flight when the process is killed abruptly. A graceful stop drains; kill -9 loses whatever was running |
Idempotency protects you against duplicates, not against loss. Parallel consumption is intended for idempotent or overwrite handlers.
retry = -1 (the annotation default) is coerced to 0 under an executor, because retrying
forever inside a task pins a concurrency permit and starves the executor. Set an explicit retry on
any parallel handler that calls a remote dependency.
Both are independently recoverable:
| Want | Use | Cost |
|---|---|---|
| At-least-once delivery |
CheckpointMode.WATERMARK (EventoBundle.Builder.setCheckpointMode / setComponentCheckpointMode) — persists the highest contiguous completed sequence; the fetch cursor stays on the in-memory dispatch frontier, so nothing is reprocessed within a run |
A crash replays the in-flight window; the dashboard's "last event" trails by it |
| Per-aggregate ordering |
ConsumerExecutors.partitioned(name, lanes) — events sharing an aggregate id pin to one lane and are applied in sequence order, one task per lane |
A hot aggregate serialises; concurrency is bounded by lanes and reduced by key skew |
partitioned needs no handler change, which is what makes parallel consumption usable by read-model
projectors that are per-aggregate sequential rather than genuinely idempotent.
Async handlers fail on their own threads, so their failures never used to reach the fetch loop: a downed dependency was met by the loop pulling at full speed and dead-lettering the entire stream. Consumers now track a transient-failure streak and back off on the same curve as a channel error.
Parallel consumers add a second term to the sizing rule. Every concurrently-executing handler that
opens a transaction (typically via a MessageHandlerInterceptor) holds a connection for its whole
task, on top of the one connection per active consumer that LockHandle already pins:
pool ≥ concurrent consumers
+ Σ(capacity of each ConsumerExecutor used by transactional handlers)
+ headroom
Cap a transactional handler's executor at its share of the pool. An executor of capacity 64 against
a 10-connection pool surfaces as acquisition timeouts — which the consumer classifies as transient
and, with the coerced retry = 0, dead-letters immediately.
Two measured notes:
-
Virtual threads are fine for transactional handlers.
VirtualThreadJdbcConcurrencyITmeasures 32 concurrentpg_sleeptransactions reaching full capacity on an 8-core box — 303 ms versus 8 s serial — so there is no carrier pinning. -
A cold Hikari pool, not the executor, throttles a freshly started async consumer. Connections
open lazily and cost more than a short handler. Raise
minimumIdleif catch-up burst latency matters.
Discovery publishes each handler's executor, so the GUI's Component Catalog marks parallel handlers
and Cluster Status → Consumers shows a Parallel row (executor names, in-flight, saturation,
transient failures).
Bundles push counters to the server, exposed on /actuator/prometheus:
-
evento.consumer.executor.{capacity,in.flight,admitted,rejected,completed,failed}— taggedbundle,instance,executor -
evento.consumer.async.{in.flight,submit.timeouts,transient.failures}— taggedbundle,instance,consumer,component
rejectedis the alerting signal — the executor reporting it is the bottleneck.
Meters are removed when a node leaves, so rolling restarts do not leak a series per dead instance. Bundles with no consumer executor push nothing at all.
- The handler is idempotent, an overwrite, or you have chosen
partitionedlanes - An explicit
retryis set if the handler calls a remote dependency - The connection pool covers
consumers + Σ(transactional executor capacity) + headroom -
evento.consumer.executor.rejectedis alerted on - You have decided whether the at-most-once window is acceptable, or set
CheckpointMode.WATERMARK
- Consumer Engines — the sequential baseline
- Consumer State Store — the pool and lock mechanics behind the sizing rule
- Observability — the full meter list
Evento Framework — Copyright 2020–2026 © Gabor Galazzo. Dual-licensed under AGPL-3.0 and a commercial licence.
This wiki documents the implementation; the repository is authoritative where the two disagree. Found something out of date? Open an issue.
Getting oriented
Internals
Operations
- Server Configuration
- Throughput and Capacity
- Observability
- Security Model
- Server REST API
- Troubleshooting
Project