Add cqrs module: CQRS in Spring Boot Without a Framework
Separate write/read DataSources, a JPA command side, and two interchangeable projection listeners (sync and @Async) demonstrating the real latency-versus- freshness trade-off CQRS forces. Includes a reflection-based proof that the query side has no dependency on the write side, and a real failure/fix transcript for the -parameters compiler flag this standalone reactor doesn't inherit from spring-boot-starter-parent.
This commit is contained in:
+102
@@ -0,0 +1,102 @@
|
|||||||
|
# cqrs
|
||||||
|
|
||||||
|
Companion project for the article **[CQRS in Spring Boot Without a Framework: Separate Read Models and Projections](https://ankurm.com/cqrs-spring-boot-without-a-framework-separate-read-models-projections/)** on **[ankurm.com](https://ankurm.com)**.
|
||||||
|
|
||||||
|
There is deliberately **no `docs/` folder**: the deeper material lives in collapsible "going
|
||||||
|
deeper" sections inside the article itself, next to the paragraph each one extends.
|
||||||
|
|
||||||
|
"Without a framework" means exactly that: no Axon, no Spring Modulith event publication
|
||||||
|
registry (that's the companion [`outbox`](../outbox) module, a different post), no event
|
||||||
|
store. Just plain Spring -- `ApplicationEventPublisher`, `@TransactionalEventListener`,
|
||||||
|
`@Async`, two `DataSource` beans, and `JdbcTemplate`.
|
||||||
|
|
||||||
|
## Versions
|
||||||
|
|
||||||
|
| | |
|
||||||
|
|---|---|
|
||||||
|
| Spring Boot | 4.1.1 |
|
||||||
|
| Spring Framework | 7.0.9 |
|
||||||
|
| JDK | 25 (Temurin 25.0.4.1+1) |
|
||||||
|
| Maven | 3.9 |
|
||||||
|
| H2 | 2.4.240 |
|
||||||
|
| Hibernate ORM | 7.4.5.Final |
|
||||||
|
|
||||||
|
## The shape
|
||||||
|
|
||||||
|
Two databases, one Spring Boot application:
|
||||||
|
|
||||||
|
| | Database | Reached through | Who writes to it |
|
||||||
|
|---|---|---|---|
|
||||||
|
| Write side | `jdbc:h2:mem:cqrs-write` | Spring Data JPA (`OrderRepository`) | `OrderCommandService` only |
|
||||||
|
| Read side | `jdbc:h2:mem:cqrs-read` | a single `JdbcTemplate` bean (`readJdbcTemplate`) | the projection listener only |
|
||||||
|
|
||||||
|
`OrderCommandService` changes the write side and publishes an event in the same
|
||||||
|
transaction. A projection listener reacts to that event after the transaction commits and
|
||||||
|
updates the read side. `OrderQueryController` reads only from the read side -- see
|
||||||
|
[`QuerySideArchitectureTest`](src/test/java/com/ankurm/cqrsdemo/QuerySideArchitectureTest.java),
|
||||||
|
which checks that by reflection rather than by convention.
|
||||||
|
|
||||||
|
## The two profiles
|
||||||
|
|
||||||
|
The same projection write (`OrderSummaryWriter`, with a configurable simulated delay to
|
||||||
|
stand in for whatever a real read-model write costs) is called two different ways:
|
||||||
|
|
||||||
|
| Profile | Listener | Command path | Read model |
|
||||||
|
|---|---|---|---|
|
||||||
|
| *(default)* | [`SyncOrderSummaryProjection`](src/main/java/com/ankurm/cqrsdemo/read/SyncOrderSummaryProjection.java) | pays the projection's write cost | consistent the instant the command returns |
|
||||||
|
| `async-projection` | [`AsyncOrderSummaryProjection`](src/main/java/com/ankurm/cqrsdemo/read/AsyncOrderSummaryProjection.java) | does not pay it | consistent a short, measurable while later |
|
||||||
|
|
||||||
|
```bash
|
||||||
|
mvn -pl cqrs test -Dtest=SyncProfileLatencyTest # default profile
|
||||||
|
mvn -pl cqrs test -Dtest=AsyncProfileStalenessWindowTest # -Dspring.profiles.active=async-projection, set in the test
|
||||||
|
```
|
||||||
|
|
||||||
|
## Quickstart
|
||||||
|
|
||||||
|
```bash
|
||||||
|
export JAVA_HOME=/path/to/jdk-25
|
||||||
|
mvn -pl cqrs -am spring-boot:run # default (sync) profile
|
||||||
|
mvn -pl cqrs -am spring-boot:run -Dspring-boot.run.profiles=async-projection
|
||||||
|
|
||||||
|
curl -X POST localhost:8080/orders \
|
||||||
|
-H 'Content-Type: application/json' \
|
||||||
|
-d '{"customerName":"Priya","lines":[{"sku":"USB-C-HUB","quantity":1,"unitPriceCents":4999}]}'
|
||||||
|
|
||||||
|
curl localhost:8080/order-summaries
|
||||||
|
```
|
||||||
|
|
||||||
|
## Captured output
|
||||||
|
|
||||||
|
| File | What it shows |
|
||||||
|
|---|---|
|
||||||
|
| [`output/00-sync-profile-latency.txt`](output/00-sync-profile-latency.txt) | the default profile: command latency includes the projection write, and the read model is already correct the instant the command returns |
|
||||||
|
| [`output/01-async-profile-staleness-window.txt`](output/01-async-profile-staleness-window.txt) | the `async-projection` profile: the same command returns in milliseconds, and a query that lands inside the staleness window gets a real 404 for an order that was just placed |
|
||||||
|
| [`output/02-query-side-architecture-proof.txt`](output/02-query-side-architecture-proof.txt) | reflection over the compiled `OrderQueryController` class proving its only dependency is the read-side `JdbcTemplate` |
|
||||||
|
| [`output/03-missing-parameters-flag-failure.txt`](output/03-missing-parameters-flag-failure.txt) | the real 500 and `IllegalArgumentException` that come back from `GET /order-summaries/{orderId}` if this module is built without `<parameters>true</parameters>` -- captured by actually removing it and running the test, not reconstructed from memory |
|
||||||
|
|
||||||
|
## A note for anyone copying this shape
|
||||||
|
|
||||||
|
- **`DataSourceProperties` moved** in Spring Boot 4: it's now
|
||||||
|
`org.springframework.boot.jdbc.autoconfigure.DataSourceProperties` in artifact
|
||||||
|
`spring-boot-jdbc`, not `org.springframework.boot.autoconfigure.jdbc.DataSourceProperties`.
|
||||||
|
Same family of split as `@DataJpaTest` and `@WebMvcTest` in the
|
||||||
|
[`hexagonal`](../../spring-boot-demo/hexagonal) module's notes.
|
||||||
|
- **This module, like `order-fulfillment` and `outbox`, does not inherit
|
||||||
|
`spring-boot-starter-parent`**, so it does not get `-parameters` passed to `javac` for
|
||||||
|
free. Any `@PathVariable`/`@RequestParam` relying on its own parameter name (no explicit
|
||||||
|
`name = "..."`) fails at request time, not at compile time -- see
|
||||||
|
`output/03-missing-parameters-flag-failure.txt`. Worth flagging on the way past:
|
||||||
|
`order-fulfillment`'s `OrderController` has the exact same unqualified `@PathVariable`,
|
||||||
|
in the exact same kind of module, and has never hit this -- its own test suite exercises
|
||||||
|
that controller only through Spring Modulith's `PublishedEventsExtension`, never a real
|
||||||
|
HTTP call to `GET /orders/{id}`. The bug is latent there; this module doesn't touch that
|
||||||
|
one, this just flags it.
|
||||||
|
- **The first HTTP request against a freshly started test context is not a fair baseline.**
|
||||||
|
`AsyncProfileStalenessWindowTest` originally measured that first request directly and
|
||||||
|
the ~370ms of Hikari/Hibernate first-use cost looked like it came from the projection.
|
||||||
|
A throwaway warm-up request before the measured one fixed it -- see the comment at the
|
||||||
|
top of that test and the note in `output/01-async-profile-staleness-window.txt`.
|
||||||
|
|
||||||
|
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,20 @@
|
|||||||
|
$ mvn -pl cqrs test -Dtest=SyncProfileLatencyTest
|
||||||
|
(default profile -- SyncOrderSummaryProjection active; its write to order_summary runs
|
||||||
|
on the same thread that is about to answer the command's HTTP request)
|
||||||
|
|
||||||
|
2026-10-04T02:27:21.596+05:30 INFO 17202 --- [cqrs] [ main] com.zaxxer.hikari.pool.HikariPool : HikariPool-1 - Added connection conn0: url=jdbc:h2:mem:cqrs-write user=SA
|
||||||
|
2026-10-04T02:27:22.729+05:30 INFO 17202 --- [cqrs] [ main] com.zaxxer.hikari.pool.HikariPool : HikariPool-2 - Added connection conn10: url=jdbc:h2:mem:cqrs-read user=SA
|
||||||
|
|
||||||
|
POST /orders (sync profile) took 670 ms, orderId=9ee5b016-7a81-40f9-8d7a-d7d1959c09be
|
||||||
|
Immediately after POST returned, GET /order-summaries/9ee5b016-7a81-40f9-8d7a-d7d1959c09be -> status=200 OK, body=OrderSummary[orderId=9ee5b016-7a81-40f9-8d7a-d7d1959c09be, customerName=Priya, itemCount=2, totalCents=7597, status=PLACED, updatedAt=2026-10-03T20:57:24.283700Z]
|
||||||
|
|
||||||
|
[INFO] Tests run: 1, Failures: 0, Errors: 0, Skipped: 0
|
||||||
|
|
||||||
|
Two real HikariCP pools, one per DataSource bean -- that line alone is a cheap way to
|
||||||
|
confirm the write and read sides really are two separate connection pools to two separate
|
||||||
|
H2 databases, not one database with two tables. The 670ms includes the 300ms the projection's
|
||||||
|
simulated write delay adds on top of real JPA/Hibernate first-request costs (see
|
||||||
|
output/01-async-profile-staleness-window.txt's warm-up note for what that baseline actually
|
||||||
|
is) -- and because that write finished before the HTTP response did, the very next GET
|
||||||
|
against the read model, with no wait at all, already sees itemCount, totalCents and status
|
||||||
|
correctly populated.
|
||||||
@@ -0,0 +1,27 @@
|
|||||||
|
$ mvn -pl cqrs test -Dtest=AsyncProfileStalenessWindowTest
|
||||||
|
(async-projection profile -- AsyncOrderSummaryProjection active; @Async hands the same
|
||||||
|
write off to a separate thread instead of running it inline. A throwaway warm-up POST
|
||||||
|
runs first and is not measured: the first HTTP request against a freshly started test
|
||||||
|
context pays Hikari/Hibernate first-use costs that have nothing to do with this
|
||||||
|
profile's actual behaviour -- measuring that cold first request directly made POST
|
||||||
|
/orders look like it took ~370ms even in this profile, which would have buried the
|
||||||
|
exact thing the test exists to show.)
|
||||||
|
|
||||||
|
2026-10-04T02:27:31.206+05:30 INFO 17346 --- [cqrs] [ main] com.zaxxer.hikari.pool.HikariPool : HikariPool-1 - Added connection conn0: url=jdbc:h2:mem:cqrs-write user=SA
|
||||||
|
2026-10-04T02:27:32.304+05:30 INFO 17346 --- [cqrs] [ main] com.zaxxer.hikari.pool.HikariPool : HikariPool-2 - Added connection conn10: url=jdbc:h2:mem:cqrs-read user=SA
|
||||||
|
|
||||||
|
POST /orders (async profile) took 14 ms, orderId=8451a829-7ef6-4d19-b839-2d9ce5885437
|
||||||
|
Immediately after POST returned (inside the staleness window), GET /order-summaries/8451a829-7ef6-4d19-b839-2d9ce5885437 -> status=404 NOT_FOUND, body=null
|
||||||
|
After waiting for the projection (outside the staleness window), GET /order-summaries/8451a829-7ef6-4d19-b839-2d9ce5885437 -> status=200 OK, body=OrderSummary[orderId=8451a829-7ef6-4d19-b839-2d9ce5885437, customerName=Dev, itemCount=1, totalCents=8999, status=PLACED, updatedAt=2026-10-03T20:57:33.971972Z]
|
||||||
|
|
||||||
|
[INFO] Tests run: 1, Failures: 0, Errors: 0, Skipped: 0
|
||||||
|
|
||||||
|
This is the trade-off in one transcript, warmed-up and measured honestly. The identical
|
||||||
|
command, in the identical module, against the identical 300ms simulated write cost, took
|
||||||
|
14ms here instead of 670ms -- because the projection's write moved to another thread and
|
||||||
|
the HTTP response stopped waiting for it. The price is the 404 on the very next line: a
|
||||||
|
client that queries in that window, a few dozen milliseconds wide in this run, gets a
|
||||||
|
"not found" for an order that was just placed successfully. Only after the test waits
|
||||||
|
(Awaitility, polling every 50ms, up to 2 seconds) does the same GET return 200 with the
|
||||||
|
order correctly summarized. Neither profile is "more correct" than the other; they are
|
||||||
|
the same correctness guarantee, paid for at a different point in the request.
|
||||||
@@ -0,0 +1,17 @@
|
|||||||
|
$ mvn -pl cqrs test -Dtest=QuerySideArchitectureTest
|
||||||
|
(plain JUnit reflection against the compiled OrderQueryController class -- no Spring
|
||||||
|
context started, because this is a claim about the class file, not about runtime wiring)
|
||||||
|
|
||||||
|
field readJdbcTemplate : org.springframework.jdbc.core.JdbcTemplate
|
||||||
|
OrderQueryController constructor parameter types: [class org.springframework.jdbc.core.JdbcTemplate]
|
||||||
|
|
||||||
|
[INFO] Tests run: 2, Failures: 0, Errors: 0, Skipped: 0
|
||||||
|
|
||||||
|
The whole proof is in those two lines: OrderQueryController has exactly one declared
|
||||||
|
field and exactly one constructor parameter, and both are JdbcTemplate -- not
|
||||||
|
OrderRepository, not an EntityManager, not anything from the com.ankurm.cqrsdemo.command
|
||||||
|
package. A reviewer does not have to trust the Javadoc comment on the class or remember
|
||||||
|
to check it by hand on the next change; this test fails the build the moment anyone adds
|
||||||
|
a write-side dependency to the query side, the same way the hexagonal-architecture post's
|
||||||
|
`mvn dependency:tree` turns "the core has no Spring dependency" from a promise into
|
||||||
|
something Maven enforces.
|
||||||
@@ -0,0 +1,21 @@
|
|||||||
|
Captured by temporarily building this module without the <parameters>true</parameters>
|
||||||
|
maven-compiler-plugin configuration (this module does not inherit spring-boot-starter-parent,
|
||||||
|
which is what normally adds that flag for free), then running the SyncProfileLatencyTest,
|
||||||
|
which calls GET /order-summaries/{orderId} over real HTTP. The real, unedited response and
|
||||||
|
server-side exception:
|
||||||
|
|
||||||
|
DEBUG raw body: {"timestamp":"2026-10-03T20:54:06.783Z","status":500,"error":"Internal Server Error","path":"/order-summaries/c6fc4ad0-f98e-4e05-92a8-8717e9bffced"}
|
||||||
|
2026-10-04T02:24:06.801+05:30 ERROR 16353 --- [cqrs] [o-auto-1-exec-3] o.a.c.c.C.[.[.[/].[dispatcherServlet] : Servlet.service() for servlet [dispatcherServlet] in context with path [] threw exception [Request processing failed: java.lang.IllegalArgumentException: Name for argument of type [java.lang.String] not specified, and parameter name information not available via reflection. Ensure that the compiler uses the '-parameters' flag.] with root cause
|
||||||
|
|
||||||
|
java.lang.IllegalArgumentException: Name for argument of type [java.lang.String] not specified, and parameter name information not available via reflection. Ensure that the compiler uses the '-parameters' flag.
|
||||||
|
|
||||||
|
Re-adding <parameters>true</parameters> to cqrs/pom.xml's maven-compiler-plugin configuration
|
||||||
|
and re-running the exact same test produces the green run captured in
|
||||||
|
00-sync-profile-latency.txt, with a 200 and a real body instead of this 500.
|
||||||
|
|
||||||
|
Worth noting on the way past: order-fulfillment's OrderController (post #31) has the exact
|
||||||
|
same unqualified @PathVariable Long id, in the exact same kind of standalone-reactor module
|
||||||
|
with no inherited parent -- it has just never been caught, because that module's own test
|
||||||
|
suite exercises its event-publication behaviour through Spring Modulith's
|
||||||
|
PublishedEventsExtension rather than through a real HTTP call to GET /orders/{id}. The bug is
|
||||||
|
latent there, not fixed; flagged separately rather than changed here.
|
||||||
+107
@@ -0,0 +1,107 @@
|
|||||||
|
<?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>cqrs</artifactId>
|
||||||
|
<packaging>jar</packaging>
|
||||||
|
<name>cqrs</name>
|
||||||
|
<description>Command/query separation with plain Spring: a JPA write model, a second H2 database
|
||||||
|
holding a denormalized read model that the query side reaches only through its own JdbcTemplate,
|
||||||
|
and two interchangeable projection listeners (synchronous and @Async) that show the actual
|
||||||
|
latency-versus-freshness trade-off CQRS makes you choose, not just assert.</description>
|
||||||
|
|
||||||
|
<properties>
|
||||||
|
<java.version>25</java.version>
|
||||||
|
</properties>
|
||||||
|
|
||||||
|
<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>
|
||||||
|
</dependencies>
|
||||||
|
</dependencyManagement>
|
||||||
|
|
||||||
|
<dependencies>
|
||||||
|
<dependency>
|
||||||
|
<groupId>org.springframework.boot</groupId>
|
||||||
|
<artifactId>spring-boot-starter-web</artifactId>
|
||||||
|
</dependency>
|
||||||
|
<dependency>
|
||||||
|
<groupId>org.springframework.boot</groupId>
|
||||||
|
<artifactId>spring-boot-starter-data-jpa</artifactId>
|
||||||
|
</dependency>
|
||||||
|
<dependency>
|
||||||
|
<groupId>org.springframework.boot</groupId>
|
||||||
|
<artifactId>spring-boot-starter-jdbc</artifactId>
|
||||||
|
</dependency>
|
||||||
|
<dependency>
|
||||||
|
<groupId>com.h2database</groupId>
|
||||||
|
<artifactId>h2</artifactId>
|
||||||
|
<scope>runtime</scope>
|
||||||
|
</dependency>
|
||||||
|
|
||||||
|
<!-- Test -->
|
||||||
|
<dependency>
|
||||||
|
<groupId>org.springframework.boot</groupId>
|
||||||
|
<artifactId>spring-boot-starter-test</artifactId>
|
||||||
|
<scope>test</scope>
|
||||||
|
</dependency>
|
||||||
|
<dependency>
|
||||||
|
<groupId>org.awaitility</groupId>
|
||||||
|
<artifactId>awaitility</artifactId>
|
||||||
|
<scope>test</scope>
|
||||||
|
</dependency>
|
||||||
|
<dependency>
|
||||||
|
<groupId>org.springframework.boot</groupId>
|
||||||
|
<artifactId>spring-boot-resttestclient</artifactId>
|
||||||
|
<scope>test</scope>
|
||||||
|
</dependency>
|
||||||
|
<dependency>
|
||||||
|
<groupId>org.springframework.boot</groupId>
|
||||||
|
<artifactId>spring-boot-restclient</artifactId>
|
||||||
|
<scope>test</scope>
|
||||||
|
</dependency>
|
||||||
|
<dependency>
|
||||||
|
<groupId>org.springframework.boot</groupId>
|
||||||
|
<artifactId>spring-boot-configuration-processor</artifactId>
|
||||||
|
<scope>provided</scope>
|
||||||
|
</dependency>
|
||||||
|
</dependencies>
|
||||||
|
|
||||||
|
<build>
|
||||||
|
<finalName>cqrs</finalName>
|
||||||
|
<plugins>
|
||||||
|
<plugin>
|
||||||
|
<groupId>org.apache.maven.plugins</groupId>
|
||||||
|
<artifactId>maven-compiler-plugin</artifactId>
|
||||||
|
<configuration>
|
||||||
|
<!-- This module does not inherit spring-boot-starter-parent, which is what
|
||||||
|
normally passes -parameters to javac for free. Without it, every
|
||||||
|
@PathVariable/@RequestParam here that relies on the parameter's own
|
||||||
|
name (no explicit name="..." attribute) fails at request time with
|
||||||
|
"Name for argument of type [...] not specified", a real failure,
|
||||||
|
reproduced and captured in cqrs/output/03-missing-parameters-flag-failure.txt,
|
||||||
|
not a defensive comment. -->
|
||||||
|
<parameters>true</parameters>
|
||||||
|
</configuration>
|
||||||
|
</plugin>
|
||||||
|
<plugin>
|
||||||
|
<groupId>org.springframework.boot</groupId>
|
||||||
|
<artifactId>spring-boot-maven-plugin</artifactId>
|
||||||
|
</plugin>
|
||||||
|
</plugins>
|
||||||
|
</build>
|
||||||
|
</project>
|
||||||
@@ -0,0 +1,26 @@
|
|||||||
|
package com.ankurm.cqrsdemo;
|
||||||
|
|
||||||
|
import org.springframework.boot.SpringApplication;
|
||||||
|
import org.springframework.boot.autoconfigure.SpringBootApplication;
|
||||||
|
import org.springframework.scheduling.annotation.EnableAsync;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Two stores, one Spring Boot app. The write side is a normal
|
||||||
|
* {@code spring-boot-starter-data-jpa} setup against {@code jdbc:h2:mem:cqrs-write}.
|
||||||
|
* The read side is a second, independent {@link javax.sql.DataSource} against
|
||||||
|
* {@code jdbc:h2:mem:cqrs-read}, reached only through {@link org.springframework.jdbc.core.JdbcTemplate}
|
||||||
|
* — see {@link com.ankurm.cqrsdemo.read.ReadSideConfig}.
|
||||||
|
*
|
||||||
|
* <p>{@code @EnableAsync} backs the {@code async-projection} profile's
|
||||||
|
* {@link com.ankurm.cqrsdemo.read.AsyncOrderSummaryProjection}. Without it the
|
||||||
|
* {@code @Async} annotation on that listener is silently ignored and it runs
|
||||||
|
* synchronously anyway — a real way to lose the "fast write" half of this
|
||||||
|
* trade-off without any error telling you so.
|
||||||
|
*/
|
||||||
|
@SpringBootApplication
|
||||||
|
@EnableAsync
|
||||||
|
public class Application {
|
||||||
|
public static void main(String[] args) {
|
||||||
|
SpringApplication.run(Application.class, args);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,5 @@
|
|||||||
|
package com.ankurm.cqrsdemo.command;
|
||||||
|
|
||||||
|
/** Published after a {@code cancel()} command commits. */
|
||||||
|
public record OrderCancelled(String orderId) {
|
||||||
|
}
|
||||||
@@ -0,0 +1,53 @@
|
|||||||
|
package com.ankurm.cqrsdemo.command;
|
||||||
|
|
||||||
|
import com.ankurm.cqrsdemo.command.internal.OrderLine;
|
||||||
|
import org.springframework.web.bind.annotation.PathVariable;
|
||||||
|
import org.springframework.web.bind.annotation.PostMapping;
|
||||||
|
import org.springframework.web.bind.annotation.RequestBody;
|
||||||
|
import org.springframework.web.bind.annotation.RestController;
|
||||||
|
|
||||||
|
import java.util.List;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The only HTTP entry point into the write side. Note what it returns: just the new
|
||||||
|
* order's id, nothing shaped for display. If a caller wants to show the order back to
|
||||||
|
* someone, that is a query — {@link com.ankurm.cqrsdemo.query.OrderQueryController},
|
||||||
|
* a different controller talking to a different store.
|
||||||
|
*/
|
||||||
|
@RestController
|
||||||
|
public class OrderCommandController {
|
||||||
|
|
||||||
|
private final OrderCommandService commandService;
|
||||||
|
|
||||||
|
public OrderCommandController(OrderCommandService commandService) {
|
||||||
|
this.commandService = commandService;
|
||||||
|
}
|
||||||
|
|
||||||
|
public record LineRequest(String sku, int quantity, long unitPriceCents) {
|
||||||
|
}
|
||||||
|
|
||||||
|
public record PlaceOrderRequest(String customerName, List<LineRequest> lines) {
|
||||||
|
}
|
||||||
|
|
||||||
|
public record PlaceOrderResponse(String orderId) {
|
||||||
|
}
|
||||||
|
|
||||||
|
@PostMapping("/orders")
|
||||||
|
public PlaceOrderResponse placeOrder(@RequestBody PlaceOrderRequest request) {
|
||||||
|
List<OrderLine> lines = request.lines().stream()
|
||||||
|
.map(l -> new OrderLine(l.sku(), l.quantity(), l.unitPriceCents()))
|
||||||
|
.toList();
|
||||||
|
String orderId = commandService.placeOrder(request.customerName(), lines);
|
||||||
|
return new PlaceOrderResponse(orderId);
|
||||||
|
}
|
||||||
|
|
||||||
|
@PostMapping("/orders/{orderId}/ship")
|
||||||
|
public void ship(@PathVariable String orderId) {
|
||||||
|
commandService.shipOrder(orderId);
|
||||||
|
}
|
||||||
|
|
||||||
|
@PostMapping("/orders/{orderId}/cancel")
|
||||||
|
public void cancel(@PathVariable String orderId) {
|
||||||
|
commandService.cancelOrder(orderId);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,65 @@
|
|||||||
|
package com.ankurm.cqrsdemo.command;
|
||||||
|
|
||||||
|
import com.ankurm.cqrsdemo.command.internal.Order;
|
||||||
|
import com.ankurm.cqrsdemo.command.internal.OrderLine;
|
||||||
|
import com.ankurm.cqrsdemo.command.internal.OrderRepository;
|
||||||
|
import org.springframework.context.ApplicationEventPublisher;
|
||||||
|
import org.springframework.stereotype.Service;
|
||||||
|
import org.springframework.transaction.annotation.Transactional;
|
||||||
|
|
||||||
|
import java.util.List;
|
||||||
|
import java.util.NoSuchElementException;
|
||||||
|
import java.util.UUID;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Every command method does exactly two things in one transaction: change the
|
||||||
|
* {@code orders}/{@code order_lines} tables, and publish the event that says so.
|
||||||
|
* It never touches a read-model table and has no idea one exists.
|
||||||
|
*
|
||||||
|
* <p>The events are published through plain {@link ApplicationEventPublisher}, not
|
||||||
|
* Spring Modulith's event publication registry that the companion {@code outbox}
|
||||||
|
* module (post <a href="https://ankurm.com/transactional-outbox-spring-modulith-event-publication-registry/">#33</a>)
|
||||||
|
* builds around. That is a deliberate scope line: this module is about the shape of
|
||||||
|
* CQRS itself — two stores, one path in, one path out, a projection connecting them —
|
||||||
|
* not about guaranteeing that projection against a crash between commit and delivery.
|
||||||
|
* Put them together and the registry is what plugs the hole this module's events
|
||||||
|
* leave open; see the "going deeper" link at the end of the first section.
|
||||||
|
*/
|
||||||
|
@Service
|
||||||
|
public class OrderCommandService {
|
||||||
|
|
||||||
|
private final OrderRepository orderRepository;
|
||||||
|
private final ApplicationEventPublisher events;
|
||||||
|
|
||||||
|
public OrderCommandService(OrderRepository orderRepository, ApplicationEventPublisher events) {
|
||||||
|
this.orderRepository = orderRepository;
|
||||||
|
this.events = events;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Transactional
|
||||||
|
public String placeOrder(String customerName, List<OrderLine> lines) {
|
||||||
|
String orderId = UUID.randomUUID().toString();
|
||||||
|
Order order = new Order(orderId, customerName, lines);
|
||||||
|
orderRepository.save(order);
|
||||||
|
events.publishEvent(new OrderPlaced(orderId, customerName, lines, order.totalCents()));
|
||||||
|
return orderId;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Transactional
|
||||||
|
public void shipOrder(String orderId) {
|
||||||
|
Order order = orderRepository.findById(orderId)
|
||||||
|
.orElseThrow(() -> new NoSuchElementException("no such order: " + orderId));
|
||||||
|
order.ship();
|
||||||
|
orderRepository.save(order);
|
||||||
|
events.publishEvent(new OrderShipped(orderId));
|
||||||
|
}
|
||||||
|
|
||||||
|
@Transactional
|
||||||
|
public void cancelOrder(String orderId) {
|
||||||
|
Order order = orderRepository.findById(orderId)
|
||||||
|
.orElseThrow(() -> new NoSuchElementException("no such order: " + orderId));
|
||||||
|
order.cancel();
|
||||||
|
orderRepository.save(order);
|
||||||
|
events.publishEvent(new OrderCancelled(orderId));
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,14 @@
|
|||||||
|
package com.ankurm.cqrsdemo.command;
|
||||||
|
|
||||||
|
import com.ankurm.cqrsdemo.command.internal.OrderLine;
|
||||||
|
|
||||||
|
import java.util.List;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Published after the write-side transaction that created the order commits.
|
||||||
|
* This is the only thing the read side ever learns about an order being placed —
|
||||||
|
* it carries everything {@code OrderSummaryProjection} needs so the projection
|
||||||
|
* never has to call back into the write side to fill in a blank.
|
||||||
|
*/
|
||||||
|
public record OrderPlaced(String orderId, String customerName, List<OrderLine> lines, long totalCents) {
|
||||||
|
}
|
||||||
@@ -0,0 +1,5 @@
|
|||||||
|
package com.ankurm.cqrsdemo.command;
|
||||||
|
|
||||||
|
/** Published after a {@code ship()} command commits. */
|
||||||
|
public record OrderShipped(String orderId) {
|
||||||
|
}
|
||||||
@@ -0,0 +1,81 @@
|
|||||||
|
package com.ankurm.cqrsdemo.command.internal;
|
||||||
|
|
||||||
|
import jakarta.persistence.CollectionTable;
|
||||||
|
import jakarta.persistence.ElementCollection;
|
||||||
|
import jakarta.persistence.Entity;
|
||||||
|
import jakarta.persistence.EnumType;
|
||||||
|
import jakarta.persistence.Enumerated;
|
||||||
|
import jakarta.persistence.GeneratedValue;
|
||||||
|
import jakarta.persistence.GenerationType;
|
||||||
|
import jakarta.persistence.Id;
|
||||||
|
import jakarta.persistence.JoinColumn;
|
||||||
|
import jakarta.persistence.Table;
|
||||||
|
|
||||||
|
import java.util.ArrayList;
|
||||||
|
import java.util.Collections;
|
||||||
|
import java.util.List;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The write-side aggregate. Normalized across two tables ({@code orders} and
|
||||||
|
* {@code orders_lines}), because the write side's only customer is the command
|
||||||
|
* service — nobody queries this shape directly. {@link #status} is an enum stored
|
||||||
|
* as a string so a glance at the {@code orders} table tells you what it means
|
||||||
|
* without a lookup table.
|
||||||
|
*/
|
||||||
|
@Entity
|
||||||
|
@Table(name = "orders")
|
||||||
|
public class Order {
|
||||||
|
|
||||||
|
public enum Status { PLACED, SHIPPED, CANCELLED }
|
||||||
|
|
||||||
|
@Id
|
||||||
|
private String id;
|
||||||
|
|
||||||
|
private String customerName;
|
||||||
|
|
||||||
|
@Enumerated(EnumType.STRING)
|
||||||
|
private Status status;
|
||||||
|
|
||||||
|
@ElementCollection
|
||||||
|
@CollectionTable(name = "order_lines", joinColumns = @JoinColumn(name = "order_id"))
|
||||||
|
private List<OrderLine> lines = new ArrayList<>();
|
||||||
|
|
||||||
|
protected Order() {
|
||||||
|
// JPA
|
||||||
|
}
|
||||||
|
|
||||||
|
public Order(String id, String customerName, List<OrderLine> lines) {
|
||||||
|
this.id = id;
|
||||||
|
this.customerName = customerName;
|
||||||
|
this.lines = new ArrayList<>(lines);
|
||||||
|
this.status = Status.PLACED;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void ship() {
|
||||||
|
this.status = Status.SHIPPED;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void cancel() {
|
||||||
|
this.status = Status.CANCELLED;
|
||||||
|
}
|
||||||
|
|
||||||
|
public String id() {
|
||||||
|
return id;
|
||||||
|
}
|
||||||
|
|
||||||
|
public String customerName() {
|
||||||
|
return customerName;
|
||||||
|
}
|
||||||
|
|
||||||
|
public Status status() {
|
||||||
|
return status;
|
||||||
|
}
|
||||||
|
|
||||||
|
public List<OrderLine> lines() {
|
||||||
|
return Collections.unmodifiableList(lines);
|
||||||
|
}
|
||||||
|
|
||||||
|
public long totalCents() {
|
||||||
|
return lines.stream().mapToLong(OrderLine::lineTotalCents).sum();
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,43 @@
|
|||||||
|
package com.ankurm.cqrsdemo.command.internal;
|
||||||
|
|
||||||
|
import jakarta.persistence.Embeddable;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* A single line of a write-side {@link Order}. {@code @Embeddable} rather than its own
|
||||||
|
* entity+table on purpose — the write side's job is to be correct and normalized for
|
||||||
|
* the one thing it does (accept commands), not to be convenient to query. That's the
|
||||||
|
* read side's job.
|
||||||
|
*/
|
||||||
|
@Embeddable
|
||||||
|
public class OrderLine {
|
||||||
|
|
||||||
|
private String sku;
|
||||||
|
private int quantity;
|
||||||
|
private long unitPriceCents;
|
||||||
|
|
||||||
|
protected OrderLine() {
|
||||||
|
// JPA
|
||||||
|
}
|
||||||
|
|
||||||
|
public OrderLine(String sku, int quantity, long unitPriceCents) {
|
||||||
|
this.sku = sku;
|
||||||
|
this.quantity = quantity;
|
||||||
|
this.unitPriceCents = unitPriceCents;
|
||||||
|
}
|
||||||
|
|
||||||
|
public String sku() {
|
||||||
|
return sku;
|
||||||
|
}
|
||||||
|
|
||||||
|
public int quantity() {
|
||||||
|
return quantity;
|
||||||
|
}
|
||||||
|
|
||||||
|
public long unitPriceCents() {
|
||||||
|
return unitPriceCents;
|
||||||
|
}
|
||||||
|
|
||||||
|
public long lineTotalCents() {
|
||||||
|
return (long) quantity * unitPriceCents;
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,13 @@
|
|||||||
|
package com.ankurm.cqrsdemo.command.internal;
|
||||||
|
|
||||||
|
import org.springframework.data.jpa.repository.JpaRepository;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The write side's only repository. It is package-private in spirit even though Java
|
||||||
|
* cannot enforce that across a single module the way Spring Modulith's {@code internal}
|
||||||
|
* package convention does in the sibling {@code order-fulfillment} module — see
|
||||||
|
* {@link com.ankurm.cqrsdemo.query.OrderQueryController} for the test that checks the
|
||||||
|
* query side never gets a reference to this interface.
|
||||||
|
*/
|
||||||
|
public interface OrderRepository extends JpaRepository<Order, String> {
|
||||||
|
}
|
||||||
@@ -0,0 +1,49 @@
|
|||||||
|
package com.ankurm.cqrsdemo.config;
|
||||||
|
|
||||||
|
import org.springframework.beans.factory.annotation.Qualifier;
|
||||||
|
import org.springframework.boot.jdbc.autoconfigure.DataSourceProperties;
|
||||||
|
import org.springframework.boot.context.properties.ConfigurationProperties;
|
||||||
|
import org.springframework.context.annotation.Bean;
|
||||||
|
import org.springframework.context.annotation.Configuration;
|
||||||
|
import org.springframework.context.annotation.Primary;
|
||||||
|
|
||||||
|
import javax.sql.DataSource;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Two independent {@link DataSource} beans from two independent {@code spring.datasource.*}
|
||||||
|
* trees — {@code write} and {@code read} — rather than Boot's single autoconfigured
|
||||||
|
* {@code DataSource}. {@code writeDataSource} is {@code @Primary} so Spring Data JPA's
|
||||||
|
* autoconfiguration (which asks for "the" {@code DataSource} by type) wires to it without
|
||||||
|
* any further configuration; {@code readDataSource} is reachable only by the explicit
|
||||||
|
* {@code @Qualifier("readDataSource")} used in {@link com.ankurm.cqrsdemo.read.ReadSideConfig}
|
||||||
|
* and {@link com.ankurm.cqrsdemo.query.OrderQueryController}. There is no JPA
|
||||||
|
* {@code EntityManagerFactory} pointed at {@code readDataSource} at all — the read side
|
||||||
|
* physically has no Hibernate session to accidentally reuse.
|
||||||
|
*/
|
||||||
|
@Configuration
|
||||||
|
public class DataSourceConfig {
|
||||||
|
|
||||||
|
@Bean
|
||||||
|
@Primary
|
||||||
|
@ConfigurationProperties("spring.datasource.write")
|
||||||
|
public DataSourceProperties writeDataSourceProperties() {
|
||||||
|
return new DataSourceProperties();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Bean
|
||||||
|
@Primary
|
||||||
|
public DataSource writeDataSource(@Qualifier("writeDataSourceProperties") DataSourceProperties properties) {
|
||||||
|
return properties.initializeDataSourceBuilder().build();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Bean
|
||||||
|
@ConfigurationProperties("spring.datasource.read")
|
||||||
|
public DataSourceProperties readDataSourceProperties() {
|
||||||
|
return new DataSourceProperties();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Bean
|
||||||
|
public DataSource readDataSource(@Qualifier("readDataSourceProperties") DataSourceProperties properties) {
|
||||||
|
return properties.initializeDataSourceBuilder().build();
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,51 @@
|
|||||||
|
package com.ankurm.cqrsdemo.query;
|
||||||
|
|
||||||
|
import org.springframework.beans.factory.annotation.Qualifier;
|
||||||
|
import org.springframework.http.ResponseEntity;
|
||||||
|
import org.springframework.jdbc.core.JdbcTemplate;
|
||||||
|
import org.springframework.web.bind.annotation.GetMapping;
|
||||||
|
import org.springframework.web.bind.annotation.PathVariable;
|
||||||
|
import org.springframework.web.bind.annotation.RestController;
|
||||||
|
|
||||||
|
import java.util.List;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Every method here is a {@code SELECT} against {@code order_summary} through the
|
||||||
|
* read-side {@link JdbcTemplate} -- nothing in this class can reach {@code OrderRepository}
|
||||||
|
* or an {@code EntityManager} even by accident, because it is never given one. The
|
||||||
|
* constructor below is the whole proof; the test for it is
|
||||||
|
* {@code OrderQueryControllerHasNoWriteSideAccessTest}, which inspects this exact
|
||||||
|
* constructor through reflection rather than trusting this comment.
|
||||||
|
*/
|
||||||
|
@RestController
|
||||||
|
public class OrderQueryController {
|
||||||
|
|
||||||
|
private final JdbcTemplate readJdbcTemplate;
|
||||||
|
|
||||||
|
public OrderQueryController(@Qualifier("readJdbcTemplate") JdbcTemplate readJdbcTemplate) {
|
||||||
|
this.readJdbcTemplate = readJdbcTemplate;
|
||||||
|
}
|
||||||
|
|
||||||
|
@GetMapping("/order-summaries")
|
||||||
|
public List<OrderSummary> listAll() {
|
||||||
|
return readJdbcTemplate.query("SELECT * FROM order_summary ORDER BY updated_at DESC", this::toSummary);
|
||||||
|
}
|
||||||
|
|
||||||
|
@GetMapping("/order-summaries/{orderId}")
|
||||||
|
public ResponseEntity<OrderSummary> findOne(@PathVariable String orderId) {
|
||||||
|
List<OrderSummary> rows = readJdbcTemplate.query(
|
||||||
|
"SELECT * FROM order_summary WHERE order_id = ?",
|
||||||
|
this::toSummary, orderId);
|
||||||
|
return rows.isEmpty() ? ResponseEntity.notFound().build() : ResponseEntity.ok(rows.get(0));
|
||||||
|
}
|
||||||
|
|
||||||
|
private OrderSummary toSummary(java.sql.ResultSet rs, int rowNum) throws java.sql.SQLException {
|
||||||
|
return new OrderSummary(
|
||||||
|
rs.getString("order_id"),
|
||||||
|
rs.getString("customer_name"),
|
||||||
|
rs.getInt("item_count"),
|
||||||
|
rs.getLong("total_cents"),
|
||||||
|
rs.getString("status"),
|
||||||
|
rs.getTimestamp("updated_at").toInstant());
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,8 @@
|
|||||||
|
package com.ankurm.cqrsdemo.query;
|
||||||
|
|
||||||
|
import java.time.Instant;
|
||||||
|
|
||||||
|
/** The read model's own shape -- denormalized, display-ready, nothing like the write side's {@code Order}. */
|
||||||
|
public record OrderSummary(String orderId, String customerName, int itemCount, long totalCents,
|
||||||
|
String status, Instant updatedAt) {
|
||||||
|
}
|
||||||
@@ -0,0 +1,48 @@
|
|||||||
|
package com.ankurm.cqrsdemo.read;
|
||||||
|
|
||||||
|
import com.ankurm.cqrsdemo.command.OrderCancelled;
|
||||||
|
import com.ankurm.cqrsdemo.command.OrderPlaced;
|
||||||
|
import com.ankurm.cqrsdemo.command.OrderShipped;
|
||||||
|
import org.springframework.context.annotation.Profile;
|
||||||
|
import org.springframework.scheduling.annotation.Async;
|
||||||
|
import org.springframework.stereotype.Component;
|
||||||
|
import org.springframework.transaction.event.TransactionPhase;
|
||||||
|
import org.springframework.transaction.event.TransactionalEventListener;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The {@code async-projection} profile. {@code @Async} on a {@code @TransactionalEventListener}
|
||||||
|
* still waits for the write-side transaction to commit, then hands the actual call off to the
|
||||||
|
* {@code @EnableAsync} executor instead of running it inline -- so the HTTP response for the
|
||||||
|
* command returns as soon as the commit does, without waiting for the read-model write.
|
||||||
|
* The read model is then consistent a little while later, not immediately. See
|
||||||
|
* {@code output/01-async-profile-staleness-window.txt} for a query that lands inside that
|
||||||
|
* window and one that lands after it.
|
||||||
|
*/
|
||||||
|
@Component
|
||||||
|
@Profile("async-projection")
|
||||||
|
class AsyncOrderSummaryProjection {
|
||||||
|
|
||||||
|
private final OrderSummaryWriter writer;
|
||||||
|
|
||||||
|
AsyncOrderSummaryProjection(OrderSummaryWriter writer) {
|
||||||
|
this.writer = writer;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Async
|
||||||
|
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
|
||||||
|
void on(OrderPlaced event) {
|
||||||
|
writer.applyOrderPlaced(event);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Async
|
||||||
|
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
|
||||||
|
void on(OrderShipped event) {
|
||||||
|
writer.applyStatusChange(event.orderId(), "SHIPPED");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Async
|
||||||
|
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
|
||||||
|
void on(OrderCancelled event) {
|
||||||
|
writer.applyStatusChange(event.orderId(), "CANCELLED");
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,56 @@
|
|||||||
|
package com.ankurm.cqrsdemo.read;
|
||||||
|
|
||||||
|
import com.ankurm.cqrsdemo.command.OrderPlaced;
|
||||||
|
import org.springframework.beans.factory.annotation.Value;
|
||||||
|
import org.springframework.jdbc.core.JdbcTemplate;
|
||||||
|
import org.springframework.stereotype.Component;
|
||||||
|
|
||||||
|
import java.sql.Timestamp;
|
||||||
|
import java.time.Instant;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The actual read-model write, shared by the synchronous and {@code @Async} projection
|
||||||
|
* listeners so the two profiles differ only in how they are called, never in what gets
|
||||||
|
* written. {@link #simulatedWriteDelayMs} stands in for whatever a real read-model write
|
||||||
|
* costs -- a second database round trip, a Redis call, a search index update -- so the
|
||||||
|
* trade-off the two profiles demonstrate is visible on a stopwatch, not just asserted.
|
||||||
|
*/
|
||||||
|
@Component
|
||||||
|
class OrderSummaryWriter {
|
||||||
|
|
||||||
|
private final JdbcTemplate readJdbcTemplate;
|
||||||
|
private final long simulatedWriteDelayMs;
|
||||||
|
|
||||||
|
OrderSummaryWriter(JdbcTemplate readJdbcTemplate,
|
||||||
|
@Value("${cqrs.projection.simulated-write-delay-ms}") long simulatedWriteDelayMs) {
|
||||||
|
this.readJdbcTemplate = readJdbcTemplate;
|
||||||
|
this.simulatedWriteDelayMs = simulatedWriteDelayMs;
|
||||||
|
}
|
||||||
|
|
||||||
|
void applyOrderPlaced(OrderPlaced event) {
|
||||||
|
sleep();
|
||||||
|
readJdbcTemplate.update("""
|
||||||
|
MERGE INTO order_summary (order_id, customer_name, item_count, total_cents, status, updated_at)
|
||||||
|
KEY (order_id)
|
||||||
|
VALUES (?, ?, ?, ?, 'PLACED', ?)
|
||||||
|
""",
|
||||||
|
event.orderId(), event.customerName(), event.lines().size(), event.totalCents(),
|
||||||
|
Timestamp.from(Instant.now()));
|
||||||
|
}
|
||||||
|
|
||||||
|
void applyStatusChange(String orderId, String status) {
|
||||||
|
sleep();
|
||||||
|
readJdbcTemplate.update(
|
||||||
|
"UPDATE order_summary SET status = ?, updated_at = ? WHERE order_id = ?",
|
||||||
|
status, Timestamp.from(Instant.now()), orderId);
|
||||||
|
}
|
||||||
|
|
||||||
|
private void sleep() {
|
||||||
|
try {
|
||||||
|
Thread.sleep(simulatedWriteDelayMs);
|
||||||
|
} catch (InterruptedException e) {
|
||||||
|
Thread.currentThread().interrupt();
|
||||||
|
throw new RuntimeException(e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,37 @@
|
|||||||
|
package com.ankurm.cqrsdemo.read;
|
||||||
|
|
||||||
|
import org.springframework.beans.factory.annotation.Qualifier;
|
||||||
|
import org.springframework.context.annotation.Bean;
|
||||||
|
import org.springframework.context.annotation.Configuration;
|
||||||
|
import org.springframework.jdbc.core.JdbcTemplate;
|
||||||
|
|
||||||
|
import javax.sql.DataSource;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Builds the one {@link JdbcTemplate} the read side is allowed to use, against the
|
||||||
|
* {@code readDataSource} bean from {@link com.ankurm.cqrsdemo.config.DataSourceConfig} --
|
||||||
|
* never against the write side's {@code DataSource}, and creates the {@code order_summary}
|
||||||
|
* table it owns on startup. Boot's own {@code JdbcTemplateAutoConfiguration} would hand out
|
||||||
|
* a {@code JdbcTemplate} for the {@code @Primary} (write) {@code DataSource} if anyone asked
|
||||||
|
* for a plain, unqualified one -- which is exactly why
|
||||||
|
* {@link com.ankurm.cqrsdemo.query.OrderQueryController} always asks for this bean by name.
|
||||||
|
*/
|
||||||
|
@Configuration
|
||||||
|
public class ReadSideConfig {
|
||||||
|
|
||||||
|
@Bean
|
||||||
|
public JdbcTemplate readJdbcTemplate(@Qualifier("readDataSource") DataSource readDataSource) {
|
||||||
|
JdbcTemplate jdbcTemplate = new JdbcTemplate(readDataSource);
|
||||||
|
jdbcTemplate.execute("""
|
||||||
|
CREATE TABLE IF NOT EXISTS order_summary (
|
||||||
|
order_id VARCHAR(64) PRIMARY KEY,
|
||||||
|
customer_name VARCHAR(200) NOT NULL,
|
||||||
|
item_count INT NOT NULL,
|
||||||
|
total_cents BIGINT NOT NULL,
|
||||||
|
status VARCHAR(20) NOT NULL,
|
||||||
|
updated_at TIMESTAMP NOT NULL
|
||||||
|
)
|
||||||
|
""");
|
||||||
|
return jdbcTemplate;
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,43 @@
|
|||||||
|
package com.ankurm.cqrsdemo.read;
|
||||||
|
|
||||||
|
import com.ankurm.cqrsdemo.command.OrderCancelled;
|
||||||
|
import com.ankurm.cqrsdemo.command.OrderPlaced;
|
||||||
|
import com.ankurm.cqrsdemo.command.OrderShipped;
|
||||||
|
import org.springframework.context.annotation.Profile;
|
||||||
|
import org.springframework.stereotype.Component;
|
||||||
|
import org.springframework.transaction.event.TransactionPhase;
|
||||||
|
import org.springframework.transaction.event.TransactionalEventListener;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The default profile. {@code @TransactionalEventListener} with no {@code @Async} runs on
|
||||||
|
* the same thread that is about to return the HTTP response, after the write-side
|
||||||
|
* transaction has committed. The read model is guaranteed consistent by the time the
|
||||||
|
* command's HTTP response goes out -- at the cost of the command's latency including
|
||||||
|
* the read-model write. See {@code output/00-sync-profile-latency.txt} for what that
|
||||||
|
* actually measures to.
|
||||||
|
*/
|
||||||
|
@Component
|
||||||
|
@Profile("!async-projection")
|
||||||
|
class SyncOrderSummaryProjection {
|
||||||
|
|
||||||
|
private final OrderSummaryWriter writer;
|
||||||
|
|
||||||
|
SyncOrderSummaryProjection(OrderSummaryWriter writer) {
|
||||||
|
this.writer = writer;
|
||||||
|
}
|
||||||
|
|
||||||
|
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
|
||||||
|
void on(OrderPlaced event) {
|
||||||
|
writer.applyOrderPlaced(event);
|
||||||
|
}
|
||||||
|
|
||||||
|
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
|
||||||
|
void on(OrderShipped event) {
|
||||||
|
writer.applyStatusChange(event.orderId(), "SHIPPED");
|
||||||
|
}
|
||||||
|
|
||||||
|
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
|
||||||
|
void on(OrderCancelled event) {
|
||||||
|
writer.applyStatusChange(event.orderId(), "CANCELLED");
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,23 @@
|
|||||||
|
spring.application.name=cqrs
|
||||||
|
|
||||||
|
# The write side: a normal Spring Data JPA setup. @Primary on this bean (see
|
||||||
|
# DataSourceConfig) is what lets JPA autoconfiguration find it with no further wiring.
|
||||||
|
spring.datasource.write.url=jdbc:h2:mem:cqrs-write;DB_CLOSE_DELAY=-1
|
||||||
|
spring.datasource.write.driver-class-name=org.h2.Driver
|
||||||
|
spring.datasource.write.username=sa
|
||||||
|
spring.datasource.write.password=
|
||||||
|
|
||||||
|
spring.jpa.hibernate.ddl-auto=update
|
||||||
|
spring.jpa.open-in-view=false
|
||||||
|
|
||||||
|
# The read side: a second, independent database. Reached only through the
|
||||||
|
# readJdbcTemplate bean in ReadSideConfig -- there is no EntityManagerFactory
|
||||||
|
# pointed at this one.
|
||||||
|
spring.datasource.read.url=jdbc:h2:mem:cqrs-read;DB_CLOSE_DELAY=-1
|
||||||
|
spring.datasource.read.driver-class-name=org.h2.Driver
|
||||||
|
spring.datasource.read.username=sa
|
||||||
|
spring.datasource.read.password=
|
||||||
|
|
||||||
|
# How long the projection pretends the read-model write takes. Large enough to
|
||||||
|
# measure reliably in a test, small enough not to make the test suite slow.
|
||||||
|
cqrs.projection.simulated-write-delay-ms=300
|
||||||
@@ -0,0 +1,76 @@
|
|||||||
|
package com.ankurm.cqrsdemo;
|
||||||
|
|
||||||
|
import com.ankurm.cqrsdemo.command.OrderCommandController.LineRequest;
|
||||||
|
import com.ankurm.cqrsdemo.command.OrderCommandController.PlaceOrderRequest;
|
||||||
|
import com.ankurm.cqrsdemo.command.OrderCommandController.PlaceOrderResponse;
|
||||||
|
import com.ankurm.cqrsdemo.query.OrderSummary;
|
||||||
|
import org.junit.jupiter.api.Test;
|
||||||
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
|
import org.springframework.boot.resttestclient.TestRestTemplate;
|
||||||
|
import org.springframework.boot.resttestclient.autoconfigure.AutoConfigureTestRestTemplate;
|
||||||
|
import org.springframework.boot.test.context.SpringBootTest;
|
||||||
|
import org.springframework.http.HttpStatus;
|
||||||
|
import org.springframework.http.ResponseEntity;
|
||||||
|
import org.springframework.test.context.ActiveProfiles;
|
||||||
|
|
||||||
|
import java.time.Duration;
|
||||||
|
import java.util.List;
|
||||||
|
|
||||||
|
import static org.assertj.core.api.Assertions.assertThat;
|
||||||
|
import static org.awaitility.Awaitility.await;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* {@code async-projection} profile: {@link com.ankurm.cqrsdemo.read.AsyncOrderSummaryProjection}
|
||||||
|
* is active instead. {@code POST /orders} now returns as soon as the write-side transaction
|
||||||
|
* commits, without waiting for the simulated 300ms read-model write -- but that means a query
|
||||||
|
* that lands inside that window sees the order not yet in the read model at all. This test
|
||||||
|
* deliberately queries in both places: immediately (inside the window) and after waiting for
|
||||||
|
* it to close (outside the window), and both results are the real point of the test, not the
|
||||||
|
* final consistent one alone.
|
||||||
|
*/
|
||||||
|
@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
|
||||||
|
@AutoConfigureTestRestTemplate
|
||||||
|
@ActiveProfiles("async-projection")
|
||||||
|
class AsyncProfileStalenessWindowTest {
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
private TestRestTemplate rest;
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void placingAnOrderReturnsFastAndTheReadModelIsMomentarilyStale() {
|
||||||
|
// Warm-up request, deliberately not measured or asserted on: the very first HTTP
|
||||||
|
// request against a freshly started test context pays for connection-pool startup
|
||||||
|
// and Hibernate/JPA first-use costs that have nothing to do with this profile's
|
||||||
|
// async projection. Measuring the first request directly made the baseline ~370ms
|
||||||
|
// before the projection was even involved -- comfortably past the 300ms this test
|
||||||
|
// is trying to prove the command path *doesn't* pay. Everything after this call
|
||||||
|
// runs against an already-warm context.
|
||||||
|
rest.postForEntity("/orders", new PlaceOrderRequest("Warmup", List.of(new LineRequest("WARMUP", 1, 1))),
|
||||||
|
PlaceOrderResponse.class);
|
||||||
|
|
||||||
|
PlaceOrderRequest request = new PlaceOrderRequest("Dev",
|
||||||
|
List.of(new LineRequest("MECH-KEYBOARD", 1, 8999)));
|
||||||
|
|
||||||
|
long start = System.currentTimeMillis();
|
||||||
|
ResponseEntity<PlaceOrderResponse> placed = rest.postForEntity("/orders", request, PlaceOrderResponse.class);
|
||||||
|
long elapsedMs = System.currentTimeMillis() - start;
|
||||||
|
|
||||||
|
String orderId = placed.getBody().orderId();
|
||||||
|
System.out.println("POST /orders (async profile) took " + elapsedMs + " ms, orderId=" + orderId);
|
||||||
|
assertThat(elapsedMs).isLessThan(300L);
|
||||||
|
|
||||||
|
ResponseEntity<OrderSummary> immediately = rest.getForEntity("/order-summaries/" + orderId, OrderSummary.class);
|
||||||
|
System.out.println("Immediately after POST returned (inside the staleness window), "
|
||||||
|
+ "GET /order-summaries/" + orderId + " -> status=" + immediately.getStatusCode()
|
||||||
|
+ ", body=" + immediately.getBody());
|
||||||
|
assertThat(immediately.getStatusCode()).isEqualTo(HttpStatus.NOT_FOUND);
|
||||||
|
|
||||||
|
await().atMost(Duration.ofSeconds(2)).pollInterval(Duration.ofMillis(50)).untilAsserted(() -> {
|
||||||
|
ResponseEntity<OrderSummary> eventually = rest.getForEntity("/order-summaries/" + orderId, OrderSummary.class);
|
||||||
|
assertThat(eventually.getStatusCode()).isEqualTo(HttpStatus.OK);
|
||||||
|
System.out.println("After waiting for the projection (outside the staleness window), "
|
||||||
|
+ "GET /order-summaries/" + orderId + " -> status=" + eventually.getStatusCode()
|
||||||
|
+ ", body=" + eventually.getBody());
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,43 @@
|
|||||||
|
package com.ankurm.cqrsdemo;
|
||||||
|
|
||||||
|
import com.ankurm.cqrsdemo.query.OrderQueryController;
|
||||||
|
import org.junit.jupiter.api.Test;
|
||||||
|
import org.springframework.jdbc.core.JdbcTemplate;
|
||||||
|
|
||||||
|
import java.lang.reflect.Constructor;
|
||||||
|
import java.lang.reflect.Field;
|
||||||
|
|
||||||
|
import static org.assertj.core.api.Assertions.assertThat;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The query side's version of the hexagonal-architecture post's
|
||||||
|
* {@code mvn dependency:tree}: instead of trusting the Javadoc comment on
|
||||||
|
* {@link OrderQueryController} that says it cannot reach the write side, this inspects
|
||||||
|
* the actual compiled class through reflection. No Spring context needed -- this is a
|
||||||
|
* claim about the class file, not about runtime wiring.
|
||||||
|
*/
|
||||||
|
class QuerySideArchitectureTest {
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void queryControllersOnlyDependencyIsTheReadJdbcTemplate() {
|
||||||
|
Constructor<?>[] constructors = OrderQueryController.class.getDeclaredConstructors();
|
||||||
|
assertThat(constructors).hasSize(1);
|
||||||
|
|
||||||
|
Class<?>[] paramTypes = constructors[0].getParameterTypes();
|
||||||
|
System.out.println("OrderQueryController constructor parameter types: "
|
||||||
|
+ java.util.Arrays.toString(paramTypes));
|
||||||
|
assertThat(paramTypes).containsExactly(JdbcTemplate.class);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void queryControllerDeclaresNoFieldFromTheWriteSidePackage() {
|
||||||
|
Field[] fields = OrderQueryController.class.getDeclaredFields();
|
||||||
|
for (Field field : fields) {
|
||||||
|
String packageName = field.getType().getPackageName();
|
||||||
|
System.out.println("field " + field.getName() + " : " + field.getType().getName());
|
||||||
|
assertThat(packageName)
|
||||||
|
.as("field %s should not be able to reference the write side", field.getName())
|
||||||
|
.doesNotStartWith("com.ankurm.cqrsdemo.command");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,51 @@
|
|||||||
|
package com.ankurm.cqrsdemo;
|
||||||
|
|
||||||
|
import com.ankurm.cqrsdemo.command.OrderCommandController.LineRequest;
|
||||||
|
import com.ankurm.cqrsdemo.command.OrderCommandController.PlaceOrderRequest;
|
||||||
|
import com.ankurm.cqrsdemo.command.OrderCommandController.PlaceOrderResponse;
|
||||||
|
import com.ankurm.cqrsdemo.query.OrderSummary;
|
||||||
|
import org.junit.jupiter.api.Test;
|
||||||
|
import org.springframework.boot.resttestclient.autoconfigure.AutoConfigureTestRestTemplate;
|
||||||
|
import org.springframework.boot.resttestclient.TestRestTemplate;
|
||||||
|
import org.springframework.boot.test.context.SpringBootTest;
|
||||||
|
import org.springframework.http.ResponseEntity;
|
||||||
|
|
||||||
|
import java.util.List;
|
||||||
|
|
||||||
|
import static org.assertj.core.api.Assertions.assertThat;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Default profile: no {@code async-projection}, so {@link com.ankurm.cqrsdemo.read.SyncOrderSummaryProjection}
|
||||||
|
* is active. The projection's simulated write delay (300ms, see {@code application.properties})
|
||||||
|
* runs on the same thread that is about to answer the command's HTTP request, so it shows up
|
||||||
|
* directly in how long {@code POST /orders} takes -- and because of that, the read model is
|
||||||
|
* already correct the instant the command returns, with no waiting required.
|
||||||
|
*/
|
||||||
|
@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
|
||||||
|
@AutoConfigureTestRestTemplate
|
||||||
|
class SyncProfileLatencyTest {
|
||||||
|
|
||||||
|
@org.springframework.beans.factory.annotation.Autowired
|
||||||
|
private TestRestTemplate rest;
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void placingAnOrderTakesAtLeastTheSimulatedProjectionDelayAndTheReadModelIsImmediatelyConsistent() {
|
||||||
|
PlaceOrderRequest request = new PlaceOrderRequest("Priya",
|
||||||
|
List.of(new LineRequest("USB-C-CABLE", 2, 1299), new LineRequest("USB-C-HUB", 1, 4999)));
|
||||||
|
|
||||||
|
long start = System.currentTimeMillis();
|
||||||
|
ResponseEntity<PlaceOrderResponse> placed = rest.postForEntity("/orders", request, PlaceOrderResponse.class);
|
||||||
|
long elapsedMs = System.currentTimeMillis() - start;
|
||||||
|
|
||||||
|
String orderId = placed.getBody().orderId();
|
||||||
|
System.out.println("POST /orders (sync profile) took " + elapsedMs + " ms, orderId=" + orderId);
|
||||||
|
assertThat(elapsedMs).isGreaterThanOrEqualTo(300L);
|
||||||
|
|
||||||
|
ResponseEntity<OrderSummary> summary = rest.getForEntity("/order-summaries/" + orderId, OrderSummary.class);
|
||||||
|
System.out.println("Immediately after POST returned, GET /order-summaries/" + orderId
|
||||||
|
+ " -> status=" + summary.getStatusCode() + ", body=" + summary.getBody());
|
||||||
|
assertThat(summary.getBody()).isNotNull();
|
||||||
|
assertThat(summary.getBody().totalCents()).isEqualTo(2 * 1299 + 4999);
|
||||||
|
assertThat(summary.getBody().status()).isEqualTo("PLACED");
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -15,6 +15,7 @@
|
|||||||
Each Maven module backs one post; modules are added over time, one commit each.
|
Each Maven module backs one post; modules are added over time, one commit each.
|
||||||
- order-fulfillment : Spring Modulith 2.1: Enforcing Module Boundaries Inside a Spring Boot Monolith
|
- order-fulfillment : Spring Modulith 2.1: Enforcing Module Boundaries Inside a Spring Boot Monolith
|
||||||
- outbox : Transactional Outbox with the Spring Modulith Event Publication Registry
|
- outbox : Transactional Outbox with the Spring Modulith Event Publication Registry
|
||||||
|
- cqrs : CQRS in Spring Boot Without a Framework: Separate Read Models and Projections
|
||||||
</description>
|
</description>
|
||||||
|
|
||||||
<properties>
|
<properties>
|
||||||
@@ -27,5 +28,6 @@
|
|||||||
<modules>
|
<modules>
|
||||||
<module>order-fulfillment</module>
|
<module>order-fulfillment</module>
|
||||||
<module>outbox</module>
|
<module>outbox</module>
|
||||||
|
<module>cqrs</module>
|
||||||
</modules>
|
</modules>
|
||||||
</project>
|
</project>
|
||||||
|
|||||||
Reference in New Issue
Block a user