Event-Driven Design: ApplicationEvent and Message Queues
In plain words: an event mechanism splits "a fact happened" from "reacting to that fact" into two separate pieces of code. Checkout only announces a fact — "an order was created" — and has no idea whether anyone sends an SMS, adds points or notifies the warehouse, nor should it care. That keeps the main flow short and stops one failing notification service from rolling back the whole order.
Five terms appear throughout; one line each:
- Event: a plain object describing what happened; it carries no handling logic at all
- Publisher: the business code that calls
publishEvent(...)and hands the event to the container - Listener: a method annotated
@EventListener; the container invokes it when an event matches its parameter type - Multicaster: the internal object that walks the listener list and calls them one by one — it is what actually decides synchronous versus asynchronous
- Message queue (MQ): a cross-process post office that turns the event into a stored message another process picks up
an in-process event is shouting across the office — the moment you speak, everyone in the room stops working and listens until you finish (synchronous! they sit inside your own call stack), and the building next door hears nothing (it cannot cross processes). A message queue is a group chat announcement — you send it, put the phone down, and whoever is free reads it later; someone who was offline can scroll back through history (persistence), and messages can be replayed one by one. So "should we adopt MQ?" really reduces to one question: does the next building need to know?

That animation is the single most important thread of this article: by default publishEvent is a for-loop on one thread. Section 3's experiment shows you three log lines whose thread names are identical.
After this article you should be able to answer three questions:
- Why does adding
@EventListenernot mean "asynchronous"? Which missing piece makes a side effect leave the main thread? - Why must the SMS go out only after commit? What belongs in each of
BEFORE_COMMIT / AFTER_COMMIT / AFTER_ROLLBACK / AFTER_COMPLETION? - When is an in-process event enough, and when do you truly need MQ? Once you have MQ, what guarantees "not lost" and "not duplicated" respectively?
Here is an OrderService almost every e-commerce project has written. The requirement is plain: "after a successful checkout, send an SMS, add points, issue a coupon, write an audit log, and notify the warehouse to prepare." A natural first draft looks like this:
@Servicepublic class OrderService { private final OrderRepository orderRepository; private final SmsService smsService; private final PointService pointService; private final CouponService couponService; private final AuditLogService auditLogService; private final WarehouseService warehouseService; public OrderService(OrderRepository orderRepository, SmsService smsService, PointService pointService, CouponService couponService, AuditLogService auditLogService, WarehouseService warehouseService) { this.orderRepository = orderRepository; this.smsService = smsService; this.pointService = pointService; this.couponService = couponService; this.auditLogService = auditLogService; this.warehouseService = warehouseService; } @Transactional public Order placeOrder(PlaceOrderCommand cmd) { Order order = orderRepository.save(Order.create(cmd)); // The core action is storing the order. // Everything below is "also needs doing" and is crammed in too. smsService.send(cmd.userId(), "order placed"); pointService.add(cmd.userId(), order.getAmount()); couponService.issue(cmd.userId(), order.getId()); auditLogService.record("PLACE_ORDER", order.getId()); warehouseService.notifyPrepare(order.getId()); return order; }}It "works", but the constructor takes six dependencies and the method body has five calls unrelated to creating an order. Three problems follow:
- The main flow keeps growing: every new "thing to do after checkout" means editing
OrderService, turning it into a god class. - Core and non-core are fused: if the coupon service takes 3 seconds, the user waits 3 extra seconds; if a notification service throws, the whole checkout rolls back — non-core logic topples core business.
- Muddy transaction boundaries: the SMS goes out before the transaction commits, so a later rollback cannot retract an "order placed" message.
decoupling "checkout" from "what happens after" is exactly what event-driven design solves. Core business only publishes a fact ("an order was created"); who cares and what they do is each listener's own decision.

The left half of that flow (publish → listeners → transaction phases) is the in-process world; the right half (send to broker → idempotent consumer → compensation and retry) is the cross-process one. The first half of this article covers the left, the second half the right — and the animation above zooms the very first box down to bytecode level.
Spring's event model has only three roles:
| Role | Type / annotation | Responsibility |
|---|---|---|
| Event | any POJO (Spring 4.2+ no longer requires ApplicationEvent) | describes "what happened" |
| Publisher | ApplicationEventPublisher | hands the event to the container |
| Listener | a method annotated with @EventListener | reacts to the event |
Define the event first. Java 17's record is the least-effort option — immutable by nature, a perfect "fact receipt":
package com.example.order.event;import java.math.BigDecimal;import java.time.LocalDateTime;// A plain immutable object works as an event; no base class neededpublic record OrderPlacedEvent(Long orderId, Long userId, BigDecimal amount, LocalDateTime placedAt) {}The publisher just injects ApplicationEventPublisher and calls publishEvent:
@Transactionalpublic Order placeOrder(PlaceOrderCommand cmd) { Order order = orderRepository.save(Order.create(cmd)); // Only state the fact "an order was created"; nobody cares who handles it publisher.publishEvent(new OrderPlacedEvent( order.getId(), cmd.userId(), order.getAmount(), LocalDateTime.now())); return order;}Listeners receive it with @EventListener; one event can have any number of listeners:
@Componentpublic class SmsListener { private static final Logger log = LoggerFactory.getLogger(SmsListener.class); @EventListener public void onOrderPlaced(OrderPlacedEvent event) { log.info("Sending order SMS: userId={}, orderId={}", event.userId(), event.orderId()); // smsService.send(...) }}publishEventonly broadcasts; the publisher never knows how many listeners exist- Listeners match by parameter type: an
OrderPlacedEventparameter receives only that type - Adding a new "thing to do after checkout" means writing a new listener — not a single line in
OrderServicechanges
Note: extending ApplicationEvent was the style before Spring 4.2, when events had to extend a base class and pass a source through the constructor. Today any object can be an event — a record, a plain DTO, even a String — which keeps code much cleaner.
The conclusion first: by default, publishEvent runs listeners synchronously on the current thread. That sounds counter-intuitive — it is "publish/subscribe", so why is it a synchronous call? But it is Spring's default. Verify it with a tiny experiment:
@RestControllerpublic class DemoController { private static final Logger log = LoggerFactory.getLogger(DemoController.class); private final ApplicationEventPublisher publisher; public DemoController(ApplicationEventPublisher publisher) { this.publisher = publisher; } @GetMapping("/place") public String place() { log.info("Publishing event, thread = {}", Thread.currentThread().getName()); publisher.publishEvent(new OrderPlacedEvent(1L, 100L, new BigDecimal("9.9"), LocalDateTime.now())); log.info("publishEvent returned, thread = {}", Thread.currentThread().getName()); return "ok"; }}@Componentpublic class SmsListener { private static final Logger log = LoggerFactory.getLogger(SmsListener.class); @EventListener public void onOrderPlaced(OrderPlacedEvent event) { log.info("Listener running, thread = {}", Thread.currentThread().getName()); }}Hitting /place produces this order:
Publishing event, thread = http-nio-8080-exec-1Listener running, thread = http-nio-8080-exec-1publishEvent returned, thread = http-nio-8080-exec-1All three thread names are identical, and publishEvent returned comes after the listener — proof that the listener ran synchronously on the publishing thread. Two consequences you must remember:
- Listener time equals request time: a 3-second remote call inside a listener adds 3 seconds to the request.
- Exceptions propagate upward: a listener throwing an exception bubbles up to the
publishEventcall site and rolls back the whole checkout. That is the real mechanism behind "non-core logic topples core business".
many people assume events are asynchronous by default, so they put slow work like SMS and push into listeners, and the endpoint RT explodes. Remember: a listener without @Async is just a synchronous method call.
Run that conclusion through the kernel lab — first, who exactly does a synchronous publish block?
Then add three listeners and see the multicaster's for-loop plus exception propagation:
The two labs above show the results. Spread the same flow into a single-step debugger and walk it line by line to actually see who does the shouting on which thread — click Next and keep your eyes on the thread name and the running total:
publishEvent(new OrderPlacedEvent(orderNo)); // 1 the business shouts// down: the ApplicationEventMulticaster broadcasts right here, same threadsmsListener.onOrderPlaced(e); // 2 @Order(1) SMS runs firstpointsListener.onOrderPlaced(e); // 3 @Order(2) points followauditListener.onOrderPlaced(e); // 4 @Order(3) audit closes// 5 publishEvent returns only after all three return; one throw skips the rest| thread | http-nio-8080-exec-1 |
| elapsed | 0 ms |
OrderController.payOrderService.placeOrderTo detach a listener from the publishing thread you need @EnableAsync plus @Async. But in production never use it bare — @Async defaults to SimpleAsyncTaskExecutor (a new thread per call, no reuse). Configure a real pool and handle the async exception "black hole":
@Configuration@EnableAsyncpublic class EventAsyncConfig implements AsyncConfigurer { private static final Logger log = LoggerFactory.getLogger(EventAsyncConfig.class); @Bean("eventExecutor") public Executor eventExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(4); executor.setMaxPoolSize(16); executor.setQueueCapacity(200); executor.setThreadNamePrefix("evt-"); // easy to spot in logs executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); executor.initialize(); return executor; } // Exceptions from async methods do not propagate to the caller; catch them here @Override public AsyncUncaughtExceptionHandler getAsyncUncaughtExceptionHandler() { return (ex, method, params) -> log.error("Async event failed: {}", method.getName(), ex); }}Point the listener at a pool to run asynchronously:
@Componentpublic class PointListener { @Async("eventExecutor") @EventListener public void onOrderPlaced(OrderPlacedEvent event) { pointService.add(event.userId(), event.amount()); }}Warning: an exception from an async listener does not affect the publisher's transaction. That is both a benefit and a trap — on failure there is no feedback, and without an AsyncUncaughtExceptionHandler the exception is swallowed silently: points were never added and you would not know.
Attach a real pool and publish the same event again; the thread name in the log literally moves house:
The price of going async is that nobody catches your exception. This branch shows exactly where it goes:
Recall the incident from Section 1: placeOrder sent the SMS, then the transaction rolled back, leaving the user with an "order placed" message for nothing. The root cause is that the side effect happened before commit. @TransactionalEventListener exists exactly for this: it binds a listener's execution to the transaction lifecycle.
| Phase | Fires when | Typical use |
|---|---|---|
BEFORE_COMMIT | before the transaction commits | last-chance validation |
AFTER_COMMIT | after a successful commit (default) | SMS, points, publishing messages — most common |
AFTER_ROLLBACK | after a rollback | logging failures, compensation, cache cleanup |
AFTER_COMPLETION | after the transaction ends (commit or rollback) | unconditional cleanup, releasing resources |

Spread the four phases along one timeline and it clicks: BEFORE_COMMIT still sits inside the "we can still take it back" zone, COMMIT is the life-or-death line for your data, and almost every side effect belongs after that line. Treat this animated timeline as a lookup table when you debug "did my listener run at all" — first work out which phase your publish point is in.

Compare the "incident" version with the "fixed" version:
// Incident version: plain @EventListener runs before commit@EventListenerpublic void onPlaced(OrderPlacedEvent event) { smsService.send(event.userId(), "order placed"); // if the TX rolls back, the SMS is already sent}// Fixed version: runs only after a successful commit@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)public void onCommitted(OrderPlacedEvent event) { smsService.send(event.userId(), "order placed");}@TransactionalEventListener(phase = TransactionPhase.AFTER_ROLLBACK)public void onRolledBack(OrderPlacedEvent event) { log.warn("Order {} rolled back, event discarded", event.orderId());}AFTER_COMMITguarantees the SMS is sent only after "data is durable", eliminating "SMS before data"AFTER_ROLLBACKgives you a chance to compensate or alert on failed orders- By the same token, cross-process MQ messages should be sent in
AFTER_COMMITtoo, or you get "message before data"
Trap: @TransactionalEventListener has a fallbackExecution flag that defaults to false. That means if the event is published with no active transaction (for example outside a @Transactional method), the listener never runs — with no error. When debugging, first confirm the publish point is inside a transaction. If you really need it to run without a transaction, set @TransactionalEventListener(fallbackExecution = true).
See for yourself how the event is parked and only released after commit:
With several listeners on one event, @Order controls ordering (smaller runs first), SpEL expressions filter by condition, and a listener can even return an object to publish a further event, forming a chain:
@Componentpublic class OrderEventListeners { // Risk check only for orders above 100 @Order(1) @EventListener(condition = "#event.amount > 100") public void riskCheck(OrderPlacedEvent event) { riskService.check(event.orderId()); } // Returning a non-null object makes Spring publish it as a new event (event chain) @Order(2) @EventListener public RewardEvent reward(OrderPlacedEvent event) { return new RewardEvent(event.userId(), event.amount()); }}@Order(1)puts the risk check before the rewardcondition = "#event.amount > 100"references the event argument via SpEL, avoiding anifin the body- Returning
RewardEventtriggers the next publish round, whichRewardListenercan handle — a chain reaction
Key point: ordering and conditions are reliable only for synchronous listeners on the same thread. Once you add @Async, listeners run in parallel and ordering is no longer guaranteed.
Ordering deserves a live look: run it once with @Order, then add @Async to one listener and watch the "audit first, notify second" agreement fall apart.
The other frequent follow-up: if a listener writes to the database itself, should it start its own transaction? Run REQUIRES_NEW against REQUIRED and see what happens to "the row already written" when the outer flow rolls back:
The event approach works beautifully within a single process but breaks the moment you deploy in a distributed way. With three instances:
- Other instances never receive it: an event published by instance A lives only in A's JVM; listeners in B and C never see it.
- No persistence: events are in-memory objects, lost on restart — no replay, no audit.
- No retry or acknowledgement: a throwing listener gets exactly one chance; there is no delivery confirmation or retry.
- No cross-service decoupling: when user and points services are separate applications, in-process events are useless.
To carry events across process boundaries you need a message queue (MQ): turn the "event" into a persisted message, relay it through a broker, and let each consumer process subscribe independently.

Every item on the left mirrors one failure mode listed above: in-process events are fast, need zero ops and can read your transaction context directly, but outside this JVM they are nothing at all. MQ buys persistence, retries, replay and cross-service reach, at the price of a network hop, a broker to operate, and a delivery semantics that only promises "at least once" — which hands idempotency to you.
Everything above is still someone else's buttons. To walk the journey from a shout to a group message yourself, this console is wired to the same kernel — every echo is computed by it:
run lab event sync and lab event async back to back and compare the two thread names — one is http-nio-8080-exec-1, the other evt-1. Once the name changes, every difference between sync and async grows out of that one cell.
Add the dependency, then configure the connection and reliability parameters. Do not memorise the dependency — tick it and see what the pom should and should not contain:
<?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>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>3.3.4</version> <!-- 版本由 BOM 统管,子依赖不写 version -->
<relativePath/>
</parent>
<groupId>com.example</groupId>
<artifactId>demo-service</artifactId>
<version>0.0.1-SNAPSHOT</version>
<properties>
<java.version>17</java.version>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>
</project>Once the dependency is in place, the connection and reliability parameters are the lines below — publisher-confirm-type and manual ACK are the first half of reliable delivery:
spring: rabbitmq: host: localhost port: 5672 username: guest password: guest publisher-confirm-type: correlated # enable publisher confirms publisher-returns: true # callback for unroutable messages listener: simple: acknowledge-mode: manual # manual ACK, confirm only on success retry: enabled: true max-attempts: 3 # up to 3 local retriesRabbitMQ's core idea is "an exchange routes messages to queues by rule". Know the three exchange types:
| Exchange type | Routing rule | Typical use |
|---|---|---|
| Direct | exact routing key match | point-to-point, routing by business type |
| Topic | wildcard routing key (* one word, # many) | an event bus, e.g. order.placed.# |
| Fanout | ignores the key, broadcasts to all bound queues | broadcasts, cache refresh |
Declare the exchange, queue and binding in Java Config, then send with RabbitTemplate:
@Configurationpublic class RabbitConfig { @Bean public TopicExchange orderExchange() { return ExchangeBuilder.topicExchange("order.exchange").durable(true).build(); } @Bean public Queue pointQueue() { return QueueBuilder.durable("order.point.queue").build(); } @Bean public Binding pointBinding(Queue pointQueue, TopicExchange orderExchange) { // messages with order.placed.xxx land in the points queue return BindingBuilder.bind(pointQueue).to(orderExchange).with("order.placed.#"); }}@Componentpublic class OrderMessagePublisher { private final RabbitTemplate rabbitTemplate; public OrderMessagePublisher(RabbitTemplate rabbitTemplate) { this.rabbitTemplate = rabbitTemplate; } public void publish(OrderPlacedEvent event) { rabbitTemplate.convertAndSend("order.exchange", "order.placed", event); }}Consume with @RabbitListener and manual ACK so "confirm only after success":
@Componentpublic class PointConsumer { private static final Logger log = LoggerFactory.getLogger(PointConsumer.class); @RabbitListener(queues = "order.point.queue") public void onOrderPlaced(OrderPlacedEvent event, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws IOException { try { pointService.add(event.userId(), event.amount()); channel.basicAck(deliveryTag, false); // success, remove it } catch (Exception ex) { log.error("Points failed, requeueing: orderId={}", event.orderId(), ex); channel.basicNack(deliveryTag, false, true); // failure, requeue } }}Attention: basicNack(deliveryTag, false, true) requeues forever, so a "poison message" loops endlessly. In production either cap the retries and route to a dead-letter queue, or check the retry count before acknowledging.
Before you hand-edit consumer counts, turn the two dials in this sandbox and watch backlog versus double-spend behave completely differently:
queue depth: 1,200 ↓consume rate: 2,480 msg/sdedupe hits: 6 (skipped)lag ETA: 0.5s
Then take a look at where a slow request actually stalls along the whole chain — the listener lives inside your own stack, so no latency chart will ever name it:
a dead-letter queue (DLX) is the returns shelf at a courier station — parcels that failed delivery three times get pulled aside onto one rack. They stop blocking the normal sorting line and they never vanish; you log them, find the cause, and decide whether to re-dispatch or destroy. Without that shelf, the station fills up with the same failing parcels again and again and new deliveries cannot get in at all.
Both decouple, but they target different things. Choose by workload character:
| Dimension | RabbitMQ | Kafka |
|---|---|---|
| Throughput | tens of thousands/s | hundreds of thousands to millions/s |
| Latency | low (ms) | low, slightly higher under batching |
| Ordering | per-queue ordering | per-partition ordering |
| Replay | deleted on ack, hard to replay | replayable within retention |
| Model | smart broker + simple queues | simple broker + smart consumers, log-based |
| Best for | business events, task dispatch, rich routing | log collection, stream processing, data pipelines |
One-line choice: for business events (a checkout or payment that triggers several side effects), prefer RabbitMQ; for massive logs, data pipelines and anything needing replay or stream processing, choose Kafka.
Once MQ is in, a new problem appears — "message reliability": how to make messages neither lost nor duplicated. First, the three moments a message can be lost:
| Loss moment | Scenario | Countermeasure |
|---|---|---|
| Before sending | business committed but the app died before sending | local message table: write the message in the same transaction |
| During sending | message sent but the broker did not receive it | publisher confirms + scheduled compensation |
| During consuming | consumed but crashed before ACK (redelivery) | idempotent consumer + manual ACK |
The idea is a local message table: write "the message to send" and the business data in the same database transaction, so both succeed or both fail. A scheduled task then scans pending messages and delivers them reliably to MQ.
CREATE TABLE t_local_message ( id BIGINT PRIMARY KEY AUTO_INCREMENT, biz_type VARCHAR(64) NOT NULL COMMENT 'business type, e.g. ORDER_PLACED', biz_id BIGINT NOT NULL COMMENT 'business key, for idempotency', payload JSON NOT NULL COMMENT 'message body', status TINYINT NOT NULL DEFAULT 0 COMMENT '0 pending 1 sent 2 consumed', retry_count INT NOT NULL DEFAULT 0 COMMENT 'retry count', next_retry_time DATETIME NOT NULL COMMENT 'next retry time', create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, UNIQUE KEY uk_biz (biz_type, biz_id) -- the unique key enforces idempotency);Write the order and the message in one transaction, sharing it: if the order is durable, the message is durable.
@Transactionalpublic Order placeOrder(PlaceOrderCommand cmd) { Order order = orderRepository.save(Order.create(cmd)); // Write the local message in the same transaction: "data there => message there" localMessageRepository.save(LocalMessage.pending( "ORDER_PLACED", order.getId(), toPayload(order))); return order;}A scheduled task scans status = 0 messages that are due, marks them sent on success, and increments retries on failure:
@Componentpublic class LocalMessageCompensator { private static final Logger log = LoggerFactory.getLogger(LocalMessageCompensator.class); private final LocalMessageRepository repository; private final RabbitTemplate rabbitTemplate; public LocalMessageCompensator(LocalMessageRepository repository, RabbitTemplate rabbitTemplate) { this.repository = repository; this.rabbitTemplate = rabbitTemplate; } @Scheduled(fixedDelay = 5000) public void compensate() { for (LocalMessage message : repository.findDue(100)) { try { rabbitTemplate.convertAndSend("order.exchange", message.bizType(), message.payload()); repository.markSent(message.getId()); } catch (Exception ex) { log.error("Compensation send failed: id={}", message.getId(), ex); repository.bumpRetry(message.getId(), nextRetryTime()); } } }}With "at least once" on the sending side, the consumer must be idempotent — dedupe by unique key and skip duplicates:
@Transactionalpublic void onOrderPlaced(OrderPlacedEvent event) { // Try to insert the idempotency key; a unique-key conflict means it was handled if (!dedupeRepository.tryInsert("ORDER_PLACED", event.orderId())) { log.info("Duplicate message detected, skipping: orderId={}", event.orderId()); return; } pointService.add(event.userId(), event.amount());}Key point: reliable delivery = local message table (no loss) + publisher confirms (proof it left) + scheduled compensation (retry failures) + idempotent consumer (no dupes). Only when all four are present does an "at least once" delivery become, for the business, an "exactly once" effect.
The four pieces are not a menu — they are four stations on a single assembly line, and a missing station is exactly where messages leak. Put them on one picture instead of memorising four nouns:

The left two stations keep it leaving: the local message table makes sure a crash cannot evaporate a pending message, and publisher confirms make sure every send is acknowledged. The right two keep it counting once: scheduled compensation scoops failed messages up and retries them, while the idempotent consumer shuts duplicate deliveries down at the business layer.
@TransactionalEventListener defaults to fallbackExecution = false, so it fires only when a transaction exists. Publish an event outside a transaction and the listener silently does nothing, costing you a long hunt. First confirm the publish point is wrapped by @Transactional.
a consumer without idempotency is the classic duplicate-charge scene. Under MQ's "at least once" semantics, a network hiccup redelivers the same message — with no idempotency key, points are added twice and stock deducted twice. The idempotency key must be a business unique key (such as orderId), never the message id.
By now your head is full of correct-sounding statements — but what actually keeps you out of trouble is spotting the intuitions that sound right and are wrong. Play a round of myth-busting: six intuitions on the left, their truths on the right, and a wrong match shows you the incident scene:
Every "symptom" row below can be copy-pasted straight into a search engine. Beginners get stuck in five places: the listener never runs, the main flow gets dragged down, @Async silently does nothing, messages are consumed twice, and messages pile up.
| Symptom (excerpt) | Real cause | 30-second self-rescue | Deep dive |
|---|---|---|---|
@TransactionalEventListener prints not a single log line, and no error either | The publish point has no active transaction, and fallbackExecution defaults to false, so the container skips the listener outright | Confirm a @Transactional sits on the call chain; if you genuinely want it to run without one, write @TransactionalEventListener(fallbackExecution = true). Do not "probe" with plain @EventListener — that reintroduces the SMS-before-commit bug | Section 5 · the phase lab above |
| Checkout intermittently returns 500, the top of the stack is a third-party call inside a listener, and the whole order rolled back | Listeners run synchronously on your thread and inside your transaction, so the exception rethrows into publishEvent | Move non-core side effects to AFTER_COMMIT and catch them inside the method body; only add a pool when you truly need parallelism | Sections 3 and 5 · #17 @Transactional |
I added @Async but the log still shows http-nio-8080-exec-1 | Three usual suspects: no @EnableAsync on a config class; publisher and listener live in the same bean (self-invocation bypasses the proxy); the async method declares a return type other than void | Check the startup log for a TaskExecutor; move the listener into its own @Component; async event listeners must return void | #40 @Async and pools |
org.springframework.core.task.SimpleAsyncTaskExecutor appears in the stack while the thread count climbs into the thousands | It is the fallback when @Async finds no candidate pool — one new thread per task, no reuse, no bound | Name a pool explicitly with @Async("eventExecutor") and register a ThreadPoolTaskExecutor, or configure spring.task.execution.pool.* | #40 @Async and pools |
AsyncUncaughtExceptionHandler is never called and async exceptions vanish | That handler only applies to async methods returning void; with a Future/return value the exception is packed into the result and stays invisible unless you call get() | Keep event listeners void, or switch to CompletableFuture and attach .exceptionally(...) | Section 4 · the async(exc) lab |
Points added twice / stock deducted twice, and the MQ console shows the same message with redelivered=true | "At least once" guarantees redelivery (ACK timeout, consumer restart, rebalance). The business side did nothing about idempotency | Build a dedupe table keyed on the business unique key (orderId) with a unique index; treat an insert conflict as already-processed. Never use the message id | Section 10 · the tuning sandbox |
com.rabbitmq.client.ShutdownSignalException: channel error; reason: "NOT_FOUND - home node 'rabbit@...'", or the queue simply does not exist | @RabbitListener(queues = "...") only declares a consumer; it does not create anything. One of exchange / queue / binding is missing | Add the three @Beans (Queue, TopicExchange, Binding) or inline them with @QueueBinding | Section 8 |
Queue depth keeps climbing (messages: 41932) while consumers report idle | Each message is too slow (a synchronous remote call), prefetch is too small, or every concurrent consumer is stuck on the same slow SQL | Dump the consumer thread stack first; raise concurrentConsumers and prefetch; split the slow work into a second queue | Section 8 · the tuning sandbox |
Tons of messages land in order.point.dlx and nobody consumes them | Retries were exhausted and routing sent them to the DLX, but you never wrote a DLX consumer — poison messages now sleep there forever | Attach a consumer to the DLX that persists, alerts and allows manual replay; put max-retry and backoff into config | The Attention note in Section 8 |
The rows above can all be rescued by the error text itself. One class of incident is less kind: the line number in the stack looks blameless, while the culprit hides in an earlier line. This stack comes from a real duplicate-SMS incident — do not read the conclusion yet, point out the frame you believe is the culprit:
The day after release, support forwards three screenshots: for one single order, the customer received three order placed SMS messages. Investigation shows the ACK of the first delivery was lost on the network, so the broker redelivered the message twice. This is the stack thrown while processing the second delivery — do not read the conclusion yet, point out the frame you believe is the culprit.
Goal: prove in five minutes that synchronous, asynchronous and AFTER_COMMIT listeners genuinely differ, using thread names and ordering as evidence.
Step 1, pom.xml (only web plus transactions; no MQ needed):
<?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>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>3.3.4</version> <relativePath/> </parent> <groupId>com.example</groupId> <artifactId>event-lab</artifactId> <version>0.0.1-SNAPSHOT</version> <properties> <java.version>17</java.version> </properties> <dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</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> </dependencies> <build> <plugins> <plugin> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-maven-plugin</artifactId> </plugin> </plugins> </build></project>Step 2, src/main/resources/application.yml:
spring: datasource: url: jdbc:h2:mem:lab;DB_CLOSE_DELAY=-1 driver-class-name: org.h2.Driver sql: init: mode: always schema-locations: classpath:schema.sqllogging: level: com.example.eventlab: DEBUG pattern: console: "%d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n"Step 3, a minimal table src/main/resources/schema.sql (this is what makes the transaction real):
CREATE TABLE IF NOT EXISTS t_order ( id BIGINT PRIMARY KEY AUTO_INCREMENT, user_id BIGINT NOT NULL, amount DECIMAL(10,2) NOT NULL);Step 4, the event, the listeners and the endpoint (package com.example.eventlab):
package com.example.eventlab;import java.math.BigDecimal;public record OrderPlacedEvent(Long orderId, Long userId, BigDecimal amount) {}package com.example.eventlab;import java.util.concurrent.ThreadLocalRandom;import org.slf4j.Logger;import org.slf4j.LoggerFactory;import org.springframework.context.event.EventListener;import org.springframework.scheduling.annotation.Async;import org.springframework.stereotype.Component;import org.springframework.transaction.event.TransactionPhase;import org.springframework.transaction.event.TransactionalEventListener;@Componentpublic class LabListeners { private static final Logger log = LoggerFactory.getLogger(LabListeners.class); @EventListener public void syncListen(OrderPlacedEvent e) { log.info("[SYNC ] orderId={} thread={}", e.orderId(), Thread.currentThread().getName()); } @Async("eventExecutor") @EventListener public void asyncListen(OrderPlacedEvent e) throws InterruptedException { long cost = ThreadLocalRandom.current().nextInt(300, 800); Thread.sleep(cost); log.info("[ASYNC] orderId={} thread={} cost={}ms", e.orderId(), Thread.currentThread().getName(), cost); } @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT) public void afterCommit(OrderPlacedEvent e) { log.info("[AFTER_COMMIT] orderId={} thread={}", e.orderId(), Thread.currentThread().getName()); } @TransactionalEventListener(phase = TransactionPhase.AFTER_ROLLBACK) public void afterRollback(OrderPlacedEvent e) { log.warn("[AFTER_ROLLBACK] orderId={} discarded", e.orderId()); }}package com.example.eventlab;import java.math.BigDecimal;import java.util.Map;import org.slf4j.Logger;import org.slf4j.LoggerFactory;import org.springframework.context.ApplicationEventPublisher;import org.springframework.context.annotation.Bean;import org.springframework.jdbc.core.JdbcTemplate;import org.springframework.scheduling.annotation.Async;import org.springframework.scheduling.annotation.EnableAsync;import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;import org.springframework.transaction.annotation.Transactional;import org.springframework.web.bind.annotation.*;@RestController@EnableAsyncclass EventLabController { private static final Logger log = LoggerFactory.getLogger(EventLabController.class); private final JdbcTemplate jdbc; private final ApplicationEventPublisher publisher; EventLabController(JdbcTemplate jdbc, ApplicationEventPublisher publisher) { this.jdbc = jdbc; this.publisher = publisher; } @Bean("eventExecutor") ThreadPoolTaskExecutor eventExecutor() { ThreadPoolTaskExecutor ex = new ThreadPoolTaskExecutor(); ex.setCorePoolSize(4); ex.setMaxPoolSize(8); ex.setQueueCapacity(100); ex.setThreadNamePrefix("evt-"); return ex; } @PostMapping("/orders") @Transactional Map<String, Object> create(@RequestParam(defaultValue = "false") boolean boom) { jdbc.update("INSERT INTO t_order(user_id, amount) VALUES (?, ?)", 100L, new BigDecimal("9.90")); Long id = jdbc.queryForObject("SELECT MAX(id) FROM t_order", Long.class); log.info("about to publish thread={}", Thread.currentThread().getName()); publisher.publishEvent(new OrderPlacedEvent(id, 100L, new BigDecimal("9.90"))); if (boom) { throw new IllegalStateException("deliberate rollback"); } return Map.of("id", id, "status", "created"); }}Step 5, start it and compare two curl calls:
# happy path: both the sync listener and AFTER_COMMIT appearcurl -s -X POST "http://localhost:8080/orders"# forced rollback: the sync listener still ran (the trap!), AFTER_COMMIT is gone, AFTER_ROLLBACK shows upcurl -s -X POST "http://localhost:8080/orders?boom=true"Expected response body for the first call:
{"id":1,"status":"created"}14:32:05.112 [http-nio-8080-exec-1] INFO c.e.e.EventLabController - about to publish thread=http-nio-8080-exec-114:32:05.113 [http-nio-8080-exec-1] INFO c.e.e.LabListeners - [SYNC ] orderId=1 thread=http-nio-8080-exec-114:32:05.114 [http-nio-8080-exec-1] INFO c.e.e.EventLabController - ... (controller returns)14:32:05.118 [http-nio-8080-exec-1] INFO c.e.e.LabListeners - [AFTER_COMMIT] orderId=1 thread=http-nio-8080-exec-114:32:05.621 [evt-2] INFO c.e.e.LabListeners - [ASYNC] orderId=1 thread=evt-2 cost=509msCheck four things against that log: ① [SYNC ] shares the publisher's thread and precedes the response; ② [AFTER_COMMIT] comes after the controller line, proving it ran post-commit; ③ [ASYNC] carries thread evt-2, has the latest timestamp and finishes only after the controller already returned; ④ with boom=true, [SYNC ] still prints, [AFTER_COMMIT] disappears and [AFTER_ROLLBACK] appears — hard evidence for the trap in Section 3.
Acceptance: state which section each of those four observations maps onto; then switch the phase to BEFORE_COMMIT and explain why the SMS can still be prevented from going out.
Change one thing only, and the conclusion flips completely:
- Delete
@Transactionalfrom the controller, leave everything else. What you will observe:[SYNC ]still prints but[AFTER_COMMIT]never appears again, with no warning at all — you have reproduced the silentfallbackExecution = falsefailure. Now set@TransactionalEventListener(fallbackExecution = true)and it comes back. - Change
@Async("eventExecutor")onasyncListento bare@Async. What you will observe: the thread prefix moves fromevt-totask-(Boot'sapplicationTaskExecutor); delete the executor bean entirely and you meet the raw per-task-thread behaviour ofSimpleAsyncTaskExecutor. - Put
Thread.sleep(3000)intosyncListen. What you will observe:curljumps from tens of milliseconds to over three seconds, and withboom=truethe slow listener holds the transactional connection open the whole time — the quantified version of "listener time equals request time". - Annotate the two sync listeners with
@Order(1)and@Order(2), then add@Asyncto one. What you will observe: relative ordering stops being stable — ten requests show both orders alternating.
Tip: after variant 1, go back to the phase lab in Section 5; the two should agree exactly.
Build yourself a small "event observability" project: one OrderPlacedEvent, three delivery shapes (sync / async / AFTER_COMMIT plus a fake MQ), and an automatic health report.
Requirements:
- An
ApplicationListenercounting publishes and average handling time per event type (wrappublishEventwithSystem.nanoTime()), printing a Top 5 every minute - A
@RestControllerexposingGET /events/replay?orderId=xxxthat re-dispatches a stored event to all listeners for incident reproduction - A minimal local message table:
t_local_messageplus a@Scheduled(fixedDelay = 5000)compensator and the state machine0 pending -> 1 sent -> 2 consumed - Idempotent consumption: unique index on
(biz_type, biz_id), treatingDuplicateKeyExceptionas already processed - Zero lines of business listener code may change
Acceptance: ① break the broker (wrong address) and push 100 orders — the count of status=1 rows in t_local_message equals the order count, none lost; ② replay the same orderId five times and points increase once with "Duplicate message detected, skipping" in the log; ③ /events/replay reproduces the original listener ordering; ④ reported average handling time stays within 20% of your manually measured P99.
without looking back, name the four things the multicaster does after publishEvent (find listeners → sort → invoke one by one → return), and say at which step "async" gets inserted.
for @EventListener, @Async @EventListener and @TransactionalEventListener(AFTER_COMMIT), what thread does each run on, and does it execute before or after commit?
why must the idempotency key be a business unique key rather than the message id? Describe one concrete incident.
which of the three loss moments does the local message table cover, and what mechanism covers the other two?
name at least three signals that tell you to stop using in-process events and adopt a real MQ.
a shout is always synchronous, async needs a pool; side effects wait for commit; cross buildings, cross processes; on the wire keep a trail, on receipt dedupe.
the essence of event-driven design is "decoupling the fact from the reaction to it". In-process, use ApplicationEvent plus @TransactionalEventListener(AFTER_COMMIT) to guarantee "data first, side effects after". Across processes, use MQ to turn events into persisted messages, and guard reliability with a local message table plus idempotent consumption. In one line: core business publishes the fact only; side effects live in their own place; reliable delivery rests on the four-piece set of table, confirms, compensation and idempotency.