Documentation

Messaging

krizaka-messaging: publish through the outbox, consume with retry then DLQ, process each message once.

At-least-once messaging on RabbitMQ, written once: an event is published by one transactional line, a consumer queue is declared by one line, a failing listener is retried with back-off and then dead-lettered, and a redelivered message is processed once.

<dependency>
  <groupId>com.krizaka</groupId>
  <artifactId>krizaka-spring-boot-starter-rabbitmq</artifactId>
</dependency>

Publish an event

krizaka:
  messaging:
    producer: krizaka-users   # stamped on every event (kz-producer)
    exchanges:
      events: platform.events # defaults: krizaka.events / krizaka.dlx
      dead-letter: platform.dlx
@Transactional
public User register(NewUser input) {
  User user = users.save(input);
  events.publish("evt.user.registered", 1, new UserRegistered(user.id(), user.email()));
  return user;
}

EventPublisher.publish writes one row to your outbox, in your transaction: the event is committed or rolled back with the state change it announces, and the relay publishes it after the commit. The body is the bare event as JSON; the envelope travels in AMQP headers.

The envelope

FieldTypeRequiredDescription
kz-typestringyesThe routing key./^evt\.[a-z0-9-]+\.[a-z0-9-]+(\.v[0-9]+)?$/
kz-versionintegeryesThe contract version.≥ 1
kz-producerstringyeskrizaka.messaging.producer of the publishing service.
kz-correlation-idstringyesThe requestId of the request that produced it (krizaka-web's MDC), or a new UUID.
kz-occurred-atstringyesWhen it happened, ISO-8601 UTC.date-time

The AMQP messageId is a UUID chosen when the row is written: a relay that crashes after publishing republishes the same id, and the consumer's deduplication drops the copy. The message also carries content_type: application/json, persistent delivery and a publish timestamp.

Consume a queue

@Bean
Declarables userEvents(MessagingExchanges exchanges) {
  return KrizakaQueues.consumer(exchanges, "krizaka.notifications.user-events",
      "evt.user.registered", "evt.user.verified");
}

@RabbitListener(queues = "krizaka.notifications.user-events")
void onRegistered(UserRegistered event, @Header(AmqpHeaders.MESSAGE_ID) String messageId) {
  if (!dedup.claim("notifications.user-events", messageId)) return; // already processed
  try { welcome(event); } catch (RuntimeException e) { dedup.release("notifications.user-events", messageId); throw e; }
}

KrizakaQueues.consumer declares the events (topic) and dead-letter (direct) exchanges, the quorum queue bound to each routing key, and <queue>.dlq. The listener receives the type of its parameter — a tolerant copy of the producer's contract, never a class named in a header.

Retry, then the dead-letter queue

The kit's listener container retries a failing listener with back-off, then republishes the message to <queue>.dlq with its original headers plus x-exception-message, x-exception-stacktrace, x-original-exchange and x-original-routingKey. A rejected message is never requeued.

krizaka:
  messaging:
    retry:            # defaults
      max-attempts: 5 # the first delivery included
      initial: 500ms
      multiplier: 2
      max: 10s

Process each message once

krizaka:
  messaging:
    dedup:
      store: jdbc # or `memory` for a service without a database

MessageDedup.claim is one INSERT arbitrated by the primary key; release gives the claim back when the handler fails, so the redelivery is processed instead of dropped.

CREATE TABLE processed_messages (
    consumer     VARCHAR(255) NOT NULL,
    message_id   VARCHAR(255) NOT NULL,
    processed_at TIMESTAMPTZ  NOT NULL DEFAULT now(),
    PRIMARY KEY (consumer, message_id)
);

The outbox

The kit never creates a table: each context owns its outbox and implements OutboxStore over it. lockPendingBatch must claim its rows (FOR UPDATE SKIP LOCKED), or two instances publish the same row. The relay (krizaka.messaging.outbox.{enabled, poll-interval, batch-size, purge-interval, retention}) runs one transaction per batch, marks each published row, backs failures off and purges after the retention window. The recommended schema is in the module's README.

Conventions

WhatConventionExample
Event routing keyevt.<aggregate>.<event>evt.user.registered
Incompatible versiona new routing key .v<n> and kz-version: nevt.user.registered.v2
Consumer queue<service>.<usage>krizaka.notifications.user-events
Dead-letter queue<queue>.dlqkrizaka.notifications.user-events.dlq

A compatible change (a new optional field) keeps the routing key and the version. An incompatible one is published under evt.x.y.v2 next to evt.x.y until every consumer has moved; no shared jar of event classes, ever.

On this page