-
Notifications
You must be signed in to change notification settings - Fork 0
Consumer State Store
Consumer state — checkpoints, saga instances, dead letters, dedupe entries — lives behind five SPIs (see Consumer Engines § 2). Two implementations ship: in-memory (tests, demos) and JDBC for Postgres and MySQL.
Module: evento-consumer-state-store/evento-consumer-state-store-jdbc, package
com.evento.consumer.state.store.jdbc.
ConsumerEngineConfig::inMemory is the builder default. It wires InMemoryConsumerLock,
InMemoryConsumerStateStore, InMemorySagaStateStore, InMemoryDeadEventQueue and
InMemoryDedupeStore, plus an observer executor of
ConsumerExecutors.virtual("observer", 32).
Fine for tests and demos. It forgets every checkpoint on restart, so it is not a production option.
ConsumerEngineConfig is the record
(ConsumerProcessor processor, ConsumerStateStore stateStore, DeadEventQueue deadEventQueue); the
remaining SPIs are wired into the ConsumerProcessor. Building the JDBC variant:
private ConsumerEngineConfig jdbcEngineConfig(EventoServer eventoServer, PerformanceService ps) {
var lock = new JdbcConsumerLock(dataSource, dialect());
var stateStore = new JdbcConsumerStateStore(dataSource, dialect());
var sagaStateStore = new JdbcSagaStateStore(dataSource, dialect(), objectMapper);
var deadEventQueue = new JdbcDeadEventQueue(dataSource, dialect(), objectMapper);
var dedupeStore = new JdbcDedupeStore(dataSource, dialect());
var processor = ConsumerProcessor.builder()
.eventoServer(eventoServer)
.lock(lock)
.stateStore(stateStore)
.sagaStateStore(sagaStateStore)
.deadEventQueue(deadEventQueue)
.dedupeStore(dedupeStore)
.performanceService(ps)
.observerExecutor(ConsumerExecutors.virtual("observer", 32))
.build();
return new ConsumerEngineConfig(processor, stateStore, deadEventQueue);
}Hand it to the builder as a method reference:
EventoBundle.Builder.builder()
.setBasePackage(…)
.setBundleId("my-bundle")
.setEventoServerMessageBusConfiguration(…)
.setConsumerEngineConfigBuilder(this::jdbcEngineConfig)
.start();SqlDialect is an enum with POSTGRES and MYSQL.
var cfg = new HikariConfig();
cfg.setJdbcUrl(jdbcUrl());
cfg.setUsername(username());
cfg.setPassword(password());
cfg.setDriverClassName(driverClassName());
cfg.setMaximumPoolSize(16);
dataSource = new HikariDataSource(cfg);
FlywayMigrator.migrate(dataSource, dialect());FlywayMigrator.migrate(dataSource, dialect) applies the V1 migration for the chosen dialect from
db/migration/{postgres,mysql}/v2/V1__init_v2_consumer_state.sql.
Four tables, created by the Flyway V1 migration per dialect:
| Table | Purpose |
|---|---|
evento_v2_consumer_state |
Checkpoint + version + enabled flag + error history |
evento_v2_saga_state |
Saga state JSON(B), plus a flat associations JSON(B) column for fast lookup |
evento_v2_dead_event |
Dead-letter entries with a retry flag |
evento_v2_dedupe |
Observer dedupe entries with a sweep window |
The flat associations column is what makes saga lookup by association a single indexed
->> ? (Postgres) / JSON_EXTRACT (MySQL) rather than a scan.
ConsumerLock provides a cross-JVM exclusive zone per consumerId:
| Dialect | Mechanism |
|---|---|
| Postgres |
pg_try_advisory_lock(hashtext(consumerId)) — session-scoped |
| MySQL |
GET_LOCK(consumerId, 0) — same pattern |
Both are session-scoped, which means each held LockHandle pins one pooled Connection for its
entire lifetime. That single fact drives all the sizing below.
LockHandle is AutoCloseable and idempotent on close(). Always release it in a
try-with-resources or finally — a failed consumer cycle otherwise leaks its pinned connection.
PgDistributedLock(the server-side command lock) is a different lock. It is JVM-or-Postgres and does not consume a pool connection per held lock. With a nullDataSourceit falls back to a JVM-only semaphore, which is what makes single-JVM integration tests work.
Two terms:
pool ≥ concurrent consumers ← one pinned connection each (LockHandle)
+ Σ(capacity of each ConsumerExecutor used by
transactional handlers) ← parallel consumers only
+ headroom ← normal query/checkpoint traffic
The second term applies only if you use Parallel Consumers: every
concurrently-executing handler that opens a transaction (typically via a
MessageHandlerInterceptor) holds a connection for its whole task.
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 retry coerced to 0 under an executor, dead-letters immediately.
Under-sizing manifests as connection-acquisition timeouts, not lock errors. If you are hunting a lock bug and finding acquisition timeouts, look at pool size first.
-
Virtual threads do not pin carriers here.
VirtualThreadJdbcConcurrencyITmeasures 32 concurrentpg_sleeptransactions reaching full capacity on an 8-core box — 303 ms versus 8 s serial. -
A cold Hikari pool, not the executor, throttles a freshly started async consumer. Connections
open lazily and cost more to open than a short handler costs to run. Raise
minimumIdleif catch-up burst latency matters.
The JDBC integration tests run against Postgres and MySQL via Testcontainers and are gated behind an environment variable, so they do not require Docker on every build:
EVENTO_RUN_JDBC_IT=true JAVA_HOME=$(/usr/libexec/java_home -v 25) \
./gradlew :evento-consumer-state-store:evento-consumer-state-store-jdbc:testRoughly 50 tests — 23 scenarios × 2 dialects. See Building and Testing.
- Consumer Engines — what uses these SPIs
- Parallel Consumers § 6 — the second sizing term in context
- Troubleshooting
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