Add outbox module: Transactional Outbox with the Spring Modulith Event Publication Registry
This commit is contained in:
@@ -0,0 +1,54 @@
|
||||
# outbox
|
||||
|
||||
Companion code for *[Transactional Outbox with the Spring Modulith Event Publication
|
||||
Registry](https://ankurm.com)* on [ankurm.com](https://ankurm.com).
|
||||
|
||||
A single `billing` module demonstrating the dual-write problem, Spring Modulith's event
|
||||
publication registry acting as a real transactional outbox, externalizing events to a
|
||||
real (embedded) Kafka broker via `spring-modulith-events-kafka`, and replaying an
|
||||
incomplete publication with `IncompleteEventPublications`. Every number in the
|
||||
companion post traces to `output/`, in order.
|
||||
|
||||
## Versions
|
||||
|
||||
| Artifact | Version |
|
||||
|---|---|
|
||||
| Spring Boot | 4.1.1 |
|
||||
| Spring Modulith | 2.1.1 (GA; verified against Maven Central's `maven-metadata.xml`) |
|
||||
| spring-modulith-events-kafka | 2.1.1 |
|
||||
| spring-kafka / spring-kafka-test | 4.1.1 (managed by the Spring Boot 4.1.1 BOM; `kafka-clients` 4.2.1) |
|
||||
| JDK | 25 (LTS) |
|
||||
| H2 | runtime, in-memory, for the demo only |
|
||||
|
||||
## Quickstart
|
||||
|
||||
```
|
||||
mvn -pl outbox -am test -Dtest=DualWriteProblemTest
|
||||
mvn -pl outbox -am test -Dtest=RegistrySafetyNetTest
|
||||
mvn -pl outbox -am test -Dtest=ReplayIncompletePublicationsTest
|
||||
```
|
||||
|
||||
Run separately on purpose -- `RegistrySafetyNetTest` starts a real embedded Kafka
|
||||
broker (`@EmbeddedKafka`), which takes several seconds and is unrelated to what the
|
||||
other two tests are checking.
|
||||
|
||||
## What each class is
|
||||
|
||||
| Class | Role |
|
||||
|---|---|
|
||||
| `billing.BillingManagement` | The outbox-safe way to issue an invoice: one local transaction, one `Invoice` row, one event publication registry row |
|
||||
| `billing.NaiveBillingService` | The "before" picture -- a separate, uncoordinated `kafka.send()` call after the database write already committed. Never wired as the module's real API; only ever called directly from `DualWriteProblemTest` |
|
||||
| `billing.InvoiceIssued` | `@Externalized("invoices::#{#this.invoiceId}")` -- routes to the `invoices` Kafka topic, keyed by the invoice id |
|
||||
| `billing.FlakyAuditListener` (test-only) | Throws on its first invocation, succeeds after -- stands in for any downstream dependency being briefly unavailable, without needing to actually break a running broker mid-test |
|
||||
|
||||
## Captured output (`output/`)
|
||||
|
||||
| File | What it captures |
|
||||
|---|---|
|
||||
| `00-dual-write-problem.txt` | The invoice committed, Kafka never received either message, and a real stack trace showing `KafkaTemplate.send()` throwing *synchronously* when the broker is totally unreachable -- narrower than, and a correction to, the usual "fire-and-forget send() fails silently" claim |
|
||||
| `01-registry-safety-net.txt` | The real production path: a real embedded Kafka broker receiving the real message, and the registry moving both listeners' rows into `EVENT_PUBLICATION_ARCHIVE` (`completion-mode=ARCHIVE` is turned on for this module -- the previous post in this series left it commented out) |
|
||||
| `02-replay-incomplete-publications.txt` | A listener failing once, the registry leaving its row incomplete, and `IncompleteEventPublications.resubmitIncompletePublicationsOlderThan(Duration)` -- a real method, verified against the compiled `events-api` jar -- redelivering it successfully |
|
||||
|
||||
No `docs/` chapter directory in this repository -- the intermediate and reference-depth
|
||||
material that would normally live there is in accordion sections inside the WordPress
|
||||
post itself.
|
||||
@@ -0,0 +1,35 @@
|
||||
$ mvn -pl outbox -am test -Dtest=DualWriteProblemTest
|
||||
(NaiveBillingService: save the Invoice, then separately call kafka.send() -- no
|
||||
coordinator between the two. spring.kafka.bootstrap-servers points at localhost:19999,
|
||||
nothing listening there)
|
||||
|
||||
2026-10-04T00:52:56.429+05:30 ERROR 5271 --- [outbox] [ main] o.s.k.support.LoggingProducerListener : Exception thrown when sending a message with key='1' and payload='WIDGET-1:1999' to topic invoices-naive:
|
||||
|
||||
org.apache.kafka.common.errors.TimeoutException: Topic invoices-naive not present in metadata after 2000 ms.
|
||||
|
||||
Exception thrown synchronously from issueInvoiceTheNaiveWay(): org.springframework.kafka.KafkaException: Send failed
|
||||
Invoice count regardless: 1 (was 0). The database write already committed inside saveInvoiceOnly() before this exception was ever thrown -- the exception came from the Kafka call, several lines later, on an already-committed row.
|
||||
|
||||
2026-10-04T00:52:58.495+05:30 ERROR 5271 --- [outbox] [ main] o.s.k.support.LoggingProducerListener : Exception thrown when sending a message with key='2' and payload='WIDGET-1:2999' to topic invoices-naive:
|
||||
|
||||
org.apache.kafka.common.errors.TimeoutException: Topic invoices-naive not present in metadata after 2000 ms.
|
||||
|
||||
Invoice count after the blocked, failed send: 2 (was 1). The row committed before the send was even attempted.
|
||||
|
||||
[INFO] Tests run: 2, Failures: 0, Errors: 0, Skipped: 0
|
||||
|
||||
Two real findings here, not one. First, the obvious one: the invoice exists in the
|
||||
database in both tests, and Kafka never received either message -- that's the dual-write
|
||||
problem, reproduced with nothing faked beyond pointing Kafka at a port with no listener.
|
||||
|
||||
Second, a more specific one that contradicts the usual "fire-and-forget send() fails
|
||||
silently" framing: when the broker can't be reached to fetch topic metadata at all,
|
||||
Spring's KafkaTemplate.send() throws synchronously, wrapping the Kafka client's own
|
||||
TimeoutException in a KafkaException, even though the method's declared return type is
|
||||
a CompletableFuture. The "silent" version of this bug is real, but it is a narrower
|
||||
case than total broker unavailability -- it's what happens when the broker is reachable
|
||||
and the topic's metadata is known, but the specific produce request times out waiting
|
||||
for an acknowledgment afterward. That failure mode surfaces only through the returned
|
||||
future or a callback, never as a synchronous exception, and is not reproduced in this
|
||||
transcript (it needs a broker that accepts the connection and then stops responding
|
||||
mid-request, not one that was never there to begin with).
|
||||
@@ -0,0 +1,20 @@
|
||||
$ mvn -pl outbox -am test -Dtest=RegistrySafetyNetTest
|
||||
(real embedded Kafka broker via @EmbeddedKafka; BillingManagement.issueInvoice()
|
||||
publishes InvoiceIssued, @Externalized("invoices::#{#this.invoiceId}") routes it
|
||||
through spring-modulith-events-kafka)
|
||||
|
||||
Kafka record received -- topic=invoices key=1 value={"invoiceId":1,"sku":"GADGET-7","amountCents":4500}
|
||||
EVENT_PUBLICATION_ARCHIVE rows for InvoiceIssued: 2
|
||||
EVENT_PUBLICATION rows still outstanding for InvoiceIssued: 0
|
||||
|
||||
[INFO] Tests run: 1, Failures: 0, Errors: 0, Skipped: 0
|
||||
|
||||
The message really arrived on the real topic, with the real routing key (the invoice's
|
||||
own id, exactly as the SpEL expression after "::" in @Externalized specifies) and the
|
||||
real JSON payload. Two archive rows for InvoiceIssued, not one, because this module's
|
||||
test sources also declare FlakyAuditListener (see output/02) -- it is a second,
|
||||
independent listener on the same event, and with completion-mode=ARCHIVE configured in
|
||||
application.properties, both listeners' completed publications moved out of
|
||||
EVENT_PUBLICATION and into EVENT_PUBLICATION_ARCHIVE once they succeeded. Zero rows
|
||||
left outstanding in EVENT_PUBLICATION for this event type confirms neither listener is
|
||||
still pending.
|
||||
@@ -0,0 +1,23 @@
|
||||
$ mvn -pl outbox -am test -Dtest=ReplayIncompletePublicationsTest
|
||||
(FlakyAuditListener throws on its first invocation of InvoiceIssued, succeeds on every
|
||||
one after -- standing in for any downstream dependency that is briefly unavailable)
|
||||
|
||||
Incomplete FlakyAuditListener publications after the simulated failure: 1
|
||||
Incomplete FlakyAuditListener publications after resubmission: 0
|
||||
Invoice 1 never changed -- the retry was purely about redelivering the event, the invoice row was correct the whole time.
|
||||
|
||||
[INFO] Tests run: 1, Failures: 0, Errors: 0, Skipped: 0
|
||||
|
||||
The sequence that actually ran: BillingManagement.issueInvoice() commits the Invoice
|
||||
row and the EVENT_PUBLICATION row for FlakyAuditListener in one local transaction.
|
||||
FlakyAuditListener.on(InvoiceIssued) is then invoked asynchronously, throws on purpose,
|
||||
and the registry leaves that row with COMPLETION_DATE still null -- "incomplete" is
|
||||
exactly that: a row with no completion date, regardless of whether the listener ever
|
||||
ran at all or ran and failed. Calling
|
||||
IncompleteEventPublications.resubmitIncompletePublicationsOlderThan(Duration.ZERO) --
|
||||
a real method on a real interface, verified against the compiled events-api 2.1.1 jar,
|
||||
not assumed from documentation prose -- re-invokes the listener for every such row.
|
||||
This time failNextAttempt is already false, so it succeeds, and the row's completion
|
||||
date is set. The invoice itself was never touched by any of this: the retry is about
|
||||
redelivering a notification, not about redoing a database write that already
|
||||
succeeded the first time.
|
||||
@@ -0,0 +1,94 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<project xmlns="http://maven.apache.org/POM/4.0.0"
|
||||
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
|
||||
<modelVersion>4.0.0</modelVersion>
|
||||
|
||||
<parent>
|
||||
<groupId>com.ankurm</groupId>
|
||||
<artifactId>spring-modulith-demo</artifactId>
|
||||
<version>1.0.0</version>
|
||||
</parent>
|
||||
|
||||
<artifactId>outbox</artifactId>
|
||||
<packaging>jar</packaging>
|
||||
|
||||
<name>outbox</name>
|
||||
<description>
|
||||
Companion code for "Transactional Outbox with the Spring Modulith Event Publication
|
||||
Registry" on ankurm.com.
|
||||
</description>
|
||||
|
||||
<dependencyManagement>
|
||||
<dependencies>
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-dependencies</artifactId>
|
||||
<version>${spring-boot.version}</version>
|
||||
<type>pom</type>
|
||||
<scope>import</scope>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.modulith</groupId>
|
||||
<artifactId>spring-modulith-bom</artifactId>
|
||||
<version>${spring-modulith.version}</version>
|
||||
<type>pom</type>
|
||||
<scope>import</scope>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
</dependencyManagement>
|
||||
|
||||
<dependencies>
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.modulith</groupId>
|
||||
<artifactId>spring-modulith-starter-jdbc</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.modulith</groupId>
|
||||
<artifactId>spring-modulith-events-kafka</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-data-jpa</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.kafka</groupId>
|
||||
<artifactId>spring-kafka</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>com.h2database</groupId>
|
||||
<artifactId>h2</artifactId>
|
||||
<scope>runtime</scope>
|
||||
</dependency>
|
||||
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-test</artifactId>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.modulith</groupId>
|
||||
<artifactId>spring-modulith-starter-test</artifactId>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.kafka</groupId>
|
||||
<artifactId>spring-kafka-test</artifactId>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
|
||||
<build>
|
||||
<plugins>
|
||||
<plugin>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-maven-plugin</artifactId>
|
||||
<version>${spring-boot.version}</version>
|
||||
</plugin>
|
||||
</plugins>
|
||||
</build>
|
||||
</project>
|
||||
@@ -0,0 +1,13 @@
|
||||
package com.ankurm.outboxdemo;
|
||||
|
||||
import org.springframework.boot.SpringApplication;
|
||||
import org.springframework.boot.autoconfigure.SpringBootApplication;
|
||||
import org.springframework.modulith.Modulithic;
|
||||
|
||||
@SpringBootApplication
|
||||
@Modulithic(systemName = "Outbox Demo")
|
||||
public class Application {
|
||||
public static void main(String[] args) {
|
||||
SpringApplication.run(Application.class, args);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,35 @@
|
||||
package com.ankurm.outboxdemo.billing;
|
||||
|
||||
import com.ankurm.outboxdemo.billing.internal.Invoice;
|
||||
import com.ankurm.outboxdemo.billing.internal.InvoiceRepository;
|
||||
import org.springframework.context.ApplicationEventPublisher;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.transaction.annotation.Transactional;
|
||||
|
||||
/**
|
||||
* The outbox-safe way to issue an invoice. The {@code Invoice} row and the event
|
||||
* publication registry row for {@link InvoiceIssued} are written in the same local
|
||||
* database transaction -- that insert is what makes this an outbox, not the event
|
||||
* itself. Whatever happens to the event afterwards (externalized to Kafka, picked up
|
||||
* by an in-process listener, or nothing at all because the process dies one instruction
|
||||
* later) cannot un-happen the fact that a row recording "this needs to go out" already
|
||||
* made it to disk, atomically with the invoice itself.
|
||||
*/
|
||||
@Service
|
||||
public class BillingManagement {
|
||||
|
||||
private final InvoiceRepository invoices;
|
||||
private final ApplicationEventPublisher events;
|
||||
|
||||
public BillingManagement(InvoiceRepository invoices, ApplicationEventPublisher events) {
|
||||
this.invoices = invoices;
|
||||
this.events = events;
|
||||
}
|
||||
|
||||
@Transactional
|
||||
public Invoice issueInvoice(String sku, int amountCents) {
|
||||
var invoice = invoices.save(new Invoice(sku, amountCents, "ISSUED"));
|
||||
events.publishEvent(new InvoiceIssued(invoice.getId(), invoice.getSku(), invoice.getAmountCents()));
|
||||
return invoice;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
package com.ankurm.outboxdemo.billing;
|
||||
|
||||
import org.springframework.modulith.events.Externalized;
|
||||
|
||||
/**
|
||||
* {@code @Externalized}'s routing syntax is {@code "target::key"} -- everything before
|
||||
* {@code ::} is the Kafka topic, everything after is a SpEL expression (root object
|
||||
* {@code #this}) used as the message key. Verified against the compiled annotation
|
||||
* (org.springframework.modulith.events.Externalized, 2.1.1): it is {@code @Target(TYPE)},
|
||||
* one aliased {@code value()/target()} String attribute, nothing else -- the routing
|
||||
* syntax itself is parsed at runtime by the externalization machinery, not encoded in
|
||||
* the annotation's own shape.
|
||||
*/
|
||||
@Externalized("invoices::#{#this.invoiceId}")
|
||||
public record InvoiceIssued(Long invoiceId, String sku, int amountCents) {
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
package com.ankurm.outboxdemo.billing;
|
||||
|
||||
import com.ankurm.outboxdemo.billing.internal.Invoice;
|
||||
import com.ankurm.outboxdemo.billing.internal.InvoiceRepository;
|
||||
import org.springframework.kafka.core.KafkaTemplate;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.transaction.annotation.Transactional;
|
||||
|
||||
/**
|
||||
* The two-systems-no-coordinator version, kept here deliberately as the "before"
|
||||
* picture -- it is never wired as the module's real API, only invoked directly from
|
||||
* {@code DualWriteProblemTest}. The database write commits and returns before this
|
||||
* method ever reaches the Kafka call: there is no transaction, local or distributed,
|
||||
* spanning both. If the process dies, or the broker is unreachable, between those two
|
||||
* lines, the invoice exists and nothing downstream will ever hear about it.
|
||||
*/
|
||||
@Service
|
||||
public class NaiveBillingService {
|
||||
|
||||
private final InvoiceRepository invoices;
|
||||
private final KafkaTemplate<String, String> kafka;
|
||||
|
||||
public NaiveBillingService(InvoiceRepository invoices, KafkaTemplate<String, String> kafka) {
|
||||
this.invoices = invoices;
|
||||
this.kafka = kafka;
|
||||
}
|
||||
|
||||
@Transactional
|
||||
public Invoice saveInvoiceOnly(String sku, int amountCents) {
|
||||
return invoices.save(new Invoice(sku, amountCents, "ISSUED"));
|
||||
}
|
||||
|
||||
/**
|
||||
* Deliberately two separate steps: {@link #saveInvoiceOnly} commits and returns
|
||||
* on its own transaction boundary before this method ever calls {@code kafka.send}.
|
||||
* Fire-and-forget: the returned {@code CompletableFuture} is never inspected, which
|
||||
* is the common case at most real call sites. If the broker is unreachable, this
|
||||
* method returns successfully anyway -- the failure happens on a Kafka client
|
||||
* thread, after this method has already returned, and nothing here ever sees it.
|
||||
*/
|
||||
public Invoice issueInvoiceTheNaiveWay(String sku, int amountCents) {
|
||||
var invoice = saveInvoiceOnly(sku, amountCents);
|
||||
kafka.send("invoices-naive", invoice.getId().toString(),
|
||||
invoice.getSku() + ":" + invoice.getAmountCents());
|
||||
return invoice;
|
||||
}
|
||||
|
||||
/**
|
||||
* Same two-step shortcut, except this call site does the "responsible" thing and
|
||||
* blocks on the send future so it can react to a failure. It still doesn't help:
|
||||
* the database write already committed in {@link #saveInvoiceOnly} by the time this
|
||||
* line runs, so catching the exception here can log it, but cannot undo the invoice
|
||||
* that is already sitting in the database with nothing downstream ever told about it.
|
||||
*/
|
||||
public Invoice issueInvoiceAndBlockOnTheSend(String sku, int amountCents) throws Exception {
|
||||
var invoice = saveInvoiceOnly(sku, amountCents);
|
||||
kafka.send("invoices-naive", invoice.getId().toString(),
|
||||
invoice.getSku() + ":" + invoice.getAmountCents())
|
||||
.get(3, java.util.concurrent.TimeUnit.SECONDS);
|
||||
return invoice;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,50 @@
|
||||
package com.ankurm.outboxdemo.billing.internal;
|
||||
|
||||
import jakarta.persistence.Entity;
|
||||
import jakarta.persistence.GeneratedValue;
|
||||
import jakarta.persistence.GenerationType;
|
||||
import jakarta.persistence.Id;
|
||||
import jakarta.persistence.Table;
|
||||
|
||||
/**
|
||||
* Deliberately named Invoice, not Order -- the previous post in this series
|
||||
* (output/00-table-name-collision.txt in the order-fulfillment module) already
|
||||
* spent a transcript on why ORDER is a reserved SQL keyword. This module picks a
|
||||
* domain noun that doesn't collide, on purpose.
|
||||
*/
|
||||
@Entity
|
||||
@Table(name = "invoices")
|
||||
public class Invoice {
|
||||
|
||||
@Id
|
||||
@GeneratedValue(strategy = GenerationType.IDENTITY)
|
||||
private Long id;
|
||||
private String sku;
|
||||
private int amountCents;
|
||||
private String status;
|
||||
|
||||
protected Invoice() {
|
||||
}
|
||||
|
||||
public Invoice(String sku, int amountCents, String status) {
|
||||
this.sku = sku;
|
||||
this.amountCents = amountCents;
|
||||
this.status = status;
|
||||
}
|
||||
|
||||
public Long getId() {
|
||||
return id;
|
||||
}
|
||||
|
||||
public String getSku() {
|
||||
return sku;
|
||||
}
|
||||
|
||||
public int getAmountCents() {
|
||||
return amountCents;
|
||||
}
|
||||
|
||||
public String getStatus() {
|
||||
return status;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,6 @@
|
||||
package com.ankurm.outboxdemo.billing.internal;
|
||||
|
||||
import org.springframework.data.jpa.repository.JpaRepository;
|
||||
|
||||
public interface InvoiceRepository extends JpaRepository<Invoice, Long> {
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
spring.application.name=outbox
|
||||
|
||||
spring.datasource.url=jdbc:h2:mem:outbox;DB_CLOSE_DELAY=-1
|
||||
spring.datasource.driver-class-name=org.h2.Driver
|
||||
spring.jpa.hibernate.ddl-auto=update
|
||||
spring.jpa.open-in-view=false
|
||||
|
||||
# The event publication registry is the outbox table. JDBC-backed here (not JPA) so
|
||||
# the registry's own schema is independent of whatever ORM the business entities use.
|
||||
spring.modulith.events.jdbc.schema-initialization.enabled=true
|
||||
|
||||
# ARCHIVE, this time -- the previous post in this series left this commented out so
|
||||
# every transcript there ran against the library's bare UPDATE default. This module is
|
||||
# about the outbox actually being operated correctly, so it is turned on here instead.
|
||||
spring.modulith.events.completion-mode=ARCHIVE
|
||||
spring.modulith.events.republish-outstanding-events-on-restart=true
|
||||
|
||||
logging.level.com.ankurm.outboxdemo=INFO
|
||||
logging.level.org.springframework.modulith.events=INFO
|
||||
@@ -0,0 +1,80 @@
|
||||
package com.ankurm.outboxdemo.billing;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.assertj.core.api.Assertions.assertThatThrownBy;
|
||||
|
||||
import com.ankurm.outboxdemo.billing.internal.InvoiceRepository;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.boot.test.context.SpringBootTest;
|
||||
import org.springframework.test.context.DynamicPropertyRegistry;
|
||||
import org.springframework.test.context.DynamicPropertySource;
|
||||
|
||||
/**
|
||||
* No embedded broker here, on purpose -- spring.kafka.bootstrap-servers points at a
|
||||
* port nothing is listening on, which is the simplest honest way to reproduce "the
|
||||
* broker is unreachable" without faking anything. Two real, independent systems,
|
||||
* called one after the other, is already enough to produce the dual-write problem;
|
||||
* nothing here is simulated beyond pointing Kafka at an address with no listener.
|
||||
*/
|
||||
@SpringBootTest
|
||||
class DualWriteProblemTest {
|
||||
|
||||
@DynamicPropertySource
|
||||
static void unreachableKafka(DynamicPropertyRegistry registry) {
|
||||
registry.add("spring.kafka.bootstrap-servers", () -> "localhost:19999");
|
||||
registry.add("spring.kafka.producer.properties.max.block.ms", () -> "2000");
|
||||
}
|
||||
|
||||
@Autowired
|
||||
NaiveBillingService naiveBillingService;
|
||||
|
||||
@Autowired
|
||||
InvoiceRepository invoices;
|
||||
|
||||
@Test
|
||||
void fireAndForgetSendStillThrowsWhenTheBrokerIsTotallyUnreachable() {
|
||||
long before = invoices.count();
|
||||
|
||||
// Expected this call to return cleanly and lose the failure silently --
|
||||
// KafkaTemplate.send() returns a CompletableFuture, so a caller that ignores it
|
||||
// "shouldn't" see an exception. Running it for real shows otherwise: when the
|
||||
// broker can't even be reached to fetch topic metadata, Spring's KafkaTemplate
|
||||
// wraps that specific failure and rethrows it synchronously from send() itself
|
||||
// (see the real stack trace in the captured transcript). The database write is
|
||||
// already unaffected either way -- it's a separate, already-committed transaction
|
||||
// by the time this line runs -- but the exception itself is not silent here.
|
||||
Exception thrown = null;
|
||||
try {
|
||||
naiveBillingService.issueInvoiceTheNaiveWay("WIDGET-1", 1999);
|
||||
} catch (Exception e) {
|
||||
thrown = e;
|
||||
}
|
||||
|
||||
assertThat(thrown).isNotNull();
|
||||
assertThat(invoices.count()).isEqualTo(before + 1);
|
||||
System.out.println("Exception thrown synchronously from issueInvoiceTheNaiveWay(): "
|
||||
+ thrown.getClass().getName() + ": " + thrown.getMessage());
|
||||
System.out.println("Invoice count regardless: " + invoices.count() + " (was " + before
|
||||
+ "). The database write already committed inside saveInvoiceOnly() before this"
|
||||
+ " exception was ever thrown -- the exception came from the Kafka call, several"
|
||||
+ " lines later, on an already-committed row.");
|
||||
}
|
||||
|
||||
@Test
|
||||
void blockingOnTheFutureSurfacesTheSameFailureAndStillCannotUndoTheWrite() {
|
||||
long before = invoices.count();
|
||||
|
||||
assertThatThrownBy(() -> naiveBillingService.issueInvoiceAndBlockOnTheSend("WIDGET-1", 2999))
|
||||
.isInstanceOf(org.springframework.kafka.KafkaException.class)
|
||||
.hasCauseInstanceOf(org.apache.kafka.common.errors.TimeoutException.class);
|
||||
|
||||
// The invoice is still there. Blocking on the future told us about the Kafka
|
||||
// failure a little sooner, but by the time we could react to it,
|
||||
// saveInvoiceOnly() had already committed on its own transaction boundary --
|
||||
// there was nothing left to roll back.
|
||||
assertThat(invoices.count()).isEqualTo(before + 1);
|
||||
System.out.println("Invoice count after the blocked, failed send: " + invoices.count()
|
||||
+ " (was " + before + "). The row committed before the send was even attempted.");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,31 @@
|
||||
package com.ankurm.outboxdemo.billing;
|
||||
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import org.springframework.modulith.events.ApplicationModuleListener;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
/**
|
||||
* Test-only listener, never part of the packaged application (it lives under
|
||||
* src/test, so component scanning only ever picks it up inside a test's
|
||||
* ApplicationContext). Stands in for "the thing that can transiently fail" --
|
||||
* in a real system this is the Kafka broker being briefly unreachable, a REST
|
||||
* call to a partner timing out, or any other downstream dependency with its own
|
||||
* availability. The event publication registry doesn't know or care which;
|
||||
* it only knows whether THIS listener, for THIS event, has completed yet.
|
||||
*/
|
||||
@Component
|
||||
public class FlakyAuditListener {
|
||||
|
||||
public static final AtomicInteger attempts = new AtomicInteger(0);
|
||||
public static volatile boolean failNextAttempt = false;
|
||||
|
||||
@ApplicationModuleListener
|
||||
public void on(InvoiceIssued event) {
|
||||
int attempt = attempts.incrementAndGet();
|
||||
if (failNextAttempt) {
|
||||
failNextAttempt = false;
|
||||
throw new IllegalStateException(
|
||||
"simulated downstream failure on attempt " + attempt + " for invoice " + event.invoiceId());
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,91 @@
|
||||
package com.ankurm.outboxdemo.billing;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.awaitility.Awaitility.await;
|
||||
|
||||
import com.ankurm.outboxdemo.billing.internal.InvoiceRepository;
|
||||
import java.time.Duration;
|
||||
import java.util.Map;
|
||||
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
||||
import org.apache.kafka.clients.consumer.KafkaConsumer;
|
||||
import org.apache.kafka.common.serialization.StringDeserializer;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.boot.test.context.SpringBootTest;
|
||||
import org.springframework.jdbc.core.JdbcTemplate;
|
||||
import org.springframework.kafka.test.context.EmbeddedKafka;
|
||||
|
||||
/**
|
||||
* The real production path: BillingManagement.issueInvoice() publishes InvoiceIssued,
|
||||
* which is @Externalized, which spring-modulith-events-kafka routes to a real broker --
|
||||
* a genuine embedded Kafka instance here, not a mock. Proves two things at once: the
|
||||
* message really arrives on the topic, and the registry really marks that listener's
|
||||
* row complete (and, with completion-mode=ARCHIVE configured for this module, moves it
|
||||
* into EVENT_PUBLICATION_ARCHIVE) once it does.
|
||||
*/
|
||||
@SpringBootTest
|
||||
@EmbeddedKafka(partitions = 1, topics = "invoices", bootstrapServersProperty = "spring.kafka.bootstrap-servers")
|
||||
class RegistrySafetyNetTest {
|
||||
|
||||
@Autowired
|
||||
BillingManagement billingManagement;
|
||||
|
||||
@Autowired
|
||||
InvoiceRepository invoices;
|
||||
|
||||
@Autowired
|
||||
JdbcTemplate jdbc;
|
||||
|
||||
@Autowired
|
||||
org.springframework.core.env.Environment env;
|
||||
|
||||
KafkaConsumer<String, String> consumer;
|
||||
|
||||
@AfterEach
|
||||
void closeConsumer() {
|
||||
if (consumer != null) {
|
||||
consumer.close();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void invoiceIssuedIsReallyPublishedToKafkaAndRegistryCompletes() {
|
||||
String bootstrapServers = env.getProperty("spring.kafka.bootstrap-servers");
|
||||
consumer = new KafkaConsumer<>(Map.of(
|
||||
"bootstrap.servers", bootstrapServers,
|
||||
"group.id", "test-group",
|
||||
"auto.offset.reset", "earliest",
|
||||
"key.deserializer", StringDeserializer.class.getName(),
|
||||
"value.deserializer", StringDeserializer.class.getName()));
|
||||
consumer.subscribe(java.util.List.of("invoices"));
|
||||
|
||||
var invoice = billingManagement.issueInvoice("GADGET-7", 4500);
|
||||
assertThat(invoices.findById(invoice.getId())).isPresent();
|
||||
|
||||
var records = new java.util.ArrayList<ConsumerRecord<String, String>>();
|
||||
await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> {
|
||||
consumer.poll(Duration.ofMillis(200)).forEach(records::add);
|
||||
assertThat(records).isNotEmpty();
|
||||
});
|
||||
|
||||
ConsumerRecord<String, String> received = records.get(0);
|
||||
assertThat(received.topic()).isEqualTo("invoices");
|
||||
assertThat(received.key()).isEqualTo(invoice.getId().toString());
|
||||
System.out.println("Kafka record received -- topic=" + received.topic()
|
||||
+ " key=" + received.key() + " value=" + received.value());
|
||||
|
||||
await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> {
|
||||
Integer archived = jdbc.queryForObject(
|
||||
"select count(*) from EVENT_PUBLICATION_ARCHIVE where EVENT_TYPE like '%InvoiceIssued%'",
|
||||
Integer.class);
|
||||
System.out.println("EVENT_PUBLICATION_ARCHIVE rows for InvoiceIssued: " + archived);
|
||||
assertThat(archived).isGreaterThanOrEqualTo(1);
|
||||
});
|
||||
|
||||
Integer stillOutstanding = jdbc.queryForObject(
|
||||
"select count(*) from EVENT_PUBLICATION where EVENT_TYPE like '%InvoiceIssued%'",
|
||||
Integer.class);
|
||||
System.out.println("EVENT_PUBLICATION rows still outstanding for InvoiceIssued: " + stillOutstanding);
|
||||
}
|
||||
}
|
||||
+90
@@ -0,0 +1,90 @@
|
||||
package com.ankurm.outboxdemo.billing;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.awaitility.Awaitility.await;
|
||||
|
||||
import java.time.Duration;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.boot.test.context.SpringBootTest;
|
||||
import org.springframework.jdbc.core.JdbcTemplate;
|
||||
import org.springframework.modulith.events.IncompleteEventPublications;
|
||||
import org.springframework.test.context.DynamicPropertyRegistry;
|
||||
import org.springframework.test.context.DynamicPropertySource;
|
||||
|
||||
/**
|
||||
* FlakyAuditListener throws on its first invocation and succeeds on every one after
|
||||
* that. This is standing in for "the downstream dependency was briefly unavailable" --
|
||||
* a Kafka broker, a partner API, anything -- without needing to actually break a real
|
||||
* broker's availability mid-test. What is NOT simulated is the registry behavior: the
|
||||
* incomplete publication really is left in EVENT_PUBLICATION with no completion date
|
||||
* after the first failure, and resubmitIncompletePublicationsOlderThan really does
|
||||
* re-invoke the listener and really does mark the row complete on success. Verified
|
||||
* against the compiled IncompleteEventPublications interface (events-api 2.1.1):
|
||||
* resubmitIncompletePublicationsOlderThan(Duration) is a real method, not a guess.
|
||||
*/
|
||||
@SpringBootTest
|
||||
class ReplayIncompletePublicationsTest {
|
||||
|
||||
// This test is about the FlakyAuditListener, not Kafka -- but spring-modulith-events-kafka
|
||||
// is still on the classpath and still tries to externalize InvoiceIssued to a default
|
||||
// localhost:9092 that doesn't exist here. Trimmed down purely so the JVM doesn't sit
|
||||
// retrying metadata fetches for 30 seconds during shutdown; it does not affect what this
|
||||
// test is actually asserting.
|
||||
@DynamicPropertySource
|
||||
static void fastFailKafka(DynamicPropertyRegistry registry) {
|
||||
registry.add("spring.kafka.producer.properties.max.block.ms", () -> "500");
|
||||
}
|
||||
|
||||
@Autowired
|
||||
BillingManagement billingManagement;
|
||||
|
||||
@Autowired
|
||||
IncompleteEventPublications incompletePublications;
|
||||
|
||||
@Autowired
|
||||
JdbcTemplate jdbc;
|
||||
|
||||
@BeforeEach
|
||||
void reset() {
|
||||
FlakyAuditListener.attempts.set(0);
|
||||
FlakyAuditListener.failNextAttempt = false;
|
||||
}
|
||||
|
||||
@Test
|
||||
void failedListenerLeavesAnIncompletePublicationThatResubmissionCompletes() {
|
||||
FlakyAuditListener.failNextAttempt = true;
|
||||
|
||||
var invoice = billingManagement.issueInvoice("GIZMO-3", 999);
|
||||
|
||||
await().atMost(Duration.ofSeconds(10)).untilAsserted(() ->
|
||||
assertThat(FlakyAuditListener.attempts.get()).isEqualTo(1));
|
||||
|
||||
Integer incompleteAfterFailure = jdbc.queryForObject(
|
||||
"select count(*) from EVENT_PUBLICATION where EVENT_TYPE like '%InvoiceIssued%' "
|
||||
+ "and LISTENER_ID like '%FlakyAuditListener%' and COMPLETION_DATE is null",
|
||||
Integer.class);
|
||||
System.out.println("Incomplete FlakyAuditListener publications after the simulated failure: "
|
||||
+ incompleteAfterFailure);
|
||||
assertThat(incompleteAfterFailure).isEqualTo(1);
|
||||
|
||||
incompletePublications.resubmitIncompletePublicationsOlderThan(Duration.ZERO);
|
||||
|
||||
await().atMost(Duration.ofSeconds(10)).untilAsserted(() ->
|
||||
assertThat(FlakyAuditListener.attempts.get()).isEqualTo(2));
|
||||
|
||||
await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> {
|
||||
Integer incompleteAfterResubmit = jdbc.queryForObject(
|
||||
"select count(*) from EVENT_PUBLICATION where EVENT_TYPE like '%InvoiceIssued%' "
|
||||
+ "and LISTENER_ID like '%FlakyAuditListener%' and COMPLETION_DATE is null",
|
||||
Integer.class);
|
||||
System.out.println("Incomplete FlakyAuditListener publications after resubmission: "
|
||||
+ incompleteAfterResubmit);
|
||||
assertThat(incompleteAfterResubmit).isEqualTo(0);
|
||||
});
|
||||
|
||||
System.out.println("Invoice " + invoice.getId() + " never changed -- the retry was purely"
|
||||
+ " about redelivering the event, the invoice row was correct the whole time.");
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user