Skip to content

Consumer Engines

Gabor Galazzo edited this page Jul 25, 2026 · 1 revision

Consumer Engines

Package com.evento.application.consumer.* in evento-bundle. A consumer is anything that reads the event stream and advances a checkpoint: a projector, a saga, or an observer.


1. The engines

Class Role
ProjectorEngine Composes ConsumerProcessor + ConsumerStateStore + DeadEventQueue. Virtual-thread run loop
SagaEngine Adds SagaStateStore; per-event saga instance lookup by association
ObserverEngine At-least-once delivery plus a DedupeStore
EngineSupervisor One virtual-thread executor per engine, plus shutdown(Duration deadline) with awaitTermination → shutdownNow
ConsumerHandle The admin surface (status, dead queue, retry) — implemented by all engines
ConsumerEngineConfig The record (ConsumerProcessor processor, ConsumerStateStore stateStore, DeadEventQueue deadEventQueue) handed to the bundle builder. The remaining SPIs (lock, saga store, dedupe store) are wired into the ConsumerProcessor
DispatchContext Groups TracingAgent, the telemetry-proxy factory and MessageHandlerInterceptor — reduces constructor arity

ConsumerProcessor owns the consume loop (fetch → process → checkpoint) and is composed by the engines rather than inherited from. It holds no state beyond per-consumer in-flight counters: correctness comes from the lock plus the optimistic version on the checkpoint, not from engine-local bookkeeping.


2. The five SPIs

In evento-common, package com.evento.common.messaging.consumer.*:

SPI Responsibility
ConsumerStateStore Checkpoint read/commit with optimistic versioning, plus the enabled flag and error history
ConsumerLock Cross-JVM exclusive zone per consumerId. LockHandle is AutoCloseable
SagaStateStore Saga instance lookup by association, plus insert / update / delete
DeadEventQueue Per-consumer dead-letter queue with a retry flag
DedupeStore Observer dedupe with sweep windows
ConsumerExecutor Named, bounded execution resource for parallel consumers — ConsumerExecutors.virtual(name, maxConcurrency), pooled(name, threads), partitioned(name, lanes), unbounded(name, delegate)

Splitting v1's single abstract ConsumerStateStore class into focused interfaces is what makes them independently implementable — you can swap the lock without reimplementing checkpointing.

Every SPI has an in-memory implementation under .impl/ (zero-setup, used by tests) and a JDBC implementation for Postgres and MySQL in evento-consumer-state-store-jdbc. See Consumer State Store.


3. Sequential by default

A consumer processes one event at a time, in order, committing its checkpoint as it goes. This is the default and nothing about it changed in 2.4.0 — a bundle that does not opt in behaves exactly as before.

ConsumerLock guarantees only one instance of a given consumer runs at a time across the whole cluster, and the optimistic version on the checkpoint catches the case where a lock was lost and regained by someone else mid-cycle.

Operational note. The JDBC ConsumerLock pins a pooled Connection for the whole lifetime of the LockHandle. Size your pool for at least one connection per concurrently-running consumer. LockHandle is AutoCloseable and idempotent on close() — always release it in a try-with-resources or finally, or a failed cycle leaks its pinned connection.


4. Parallel consumption (opt-in)

An @EventHandler(executor = "name") dispatches to a named ConsumerExecutor instead of running inline, giving parallel consumption for idempotent or overwrite handlers. Executors are registered on the bundle builder, shared bundle-wide by name, and resolved per event by ConsumerExecutorResolver; an unknown name fails start-up (ConsumerExecutorValidator) rather than silently degrading to sequential.

The defining rule — the checkpoint advances when a task starts, not when it completes — and the guarantees it trades away have their own page: Parallel Consumers.


5. Interceptors

MessageHandlerInterceptor keeps working under parallel dispatch because before → handler → after/onException always run on one thread — the executor's task thread. That is what makes thread-bound transaction management viable.

Two requirements for implementors:

  • be thread-safe, and
  • keep per-invocation state in ThreadLocals.

All 24 interceptor methods are default, so you override only the hooks you need.


6. Shutdown

EngineSupervisor.shutdown(Duration deadline) does awaitTermination then shutdownNow. A graceful stop drains: supervisor → engine drain → BoundedConsumerExecutor.shutdown awaits quiescence. A kill -9 does not — see the at-most-once window in Parallel Consumers.


See also

Clone this wiki locally