事件驱动:ApplicationEvent 与消息队列集成
先说人话:事件机制就是「把『发生了一件事实』和『对这件事做出反应』拆成两拨代码」。下单这个动作只负责宣布一条事实——「订单已创建」;至于是谁要为此发短信、谁要加积分、谁要通知仓库,下单代码一概不知道、也不关心。这样主流程就不会被一堆「顺带要做的事」越拖越长,也不会因为某个通知服务挂了而整单回滚。
有五个词后面会反复出现,先各给一句话解释:
- 事件(Event):一个普通对象,只描述「发生了什么」,本身不带任何处理逻辑
- 发布者(Publisher):调用
publishEvent(...)把事件交给容器的那段业务代码 - 监听器(Listener):标了
@EventListener的方法,容器看到匹配的事件就调它 - 多播器(Multicaster):容器内部那个「拿着名单挨个打电话」的对象,真正决定同步还是异步的是它
- 消息队列(MQ):跨进程的邮局,把事件变成一条存下来的消息,由别的进程去取
进程内事件像在办公室里喊一嗓子——你一开口,整间屋子的人都停下手里的事听完才继续干活(同步!你的方法栈里就卡着他们),而且隔壁楼的人根本听不见(跨不了进程)。消息队列像微信群通知——你发完就把手机扣下,谁有空谁看、没网的人上线还能翻到聊天记录(持久化),群消息还能一条条重放(可回放)。所以「要不要上 MQ」的本质问题是:这件事需要通知隔壁楼吗?

上面这张动图是本篇最重要的一条主线:默认情况下 publishEvent 就是一根线程里的 for 循环。第七节的实验会让你亲眼看到三行日志的线程名一模一样。
学完这一篇,你应该能回答三个问题:
- 为什么「加了
@EventListener」不等于「异步」?想让副作用离开主线程,到底缺了哪一块拼图? - 短信为什么必须在事务提交后才发?
BEFORE_COMMIT / AFTER_COMMIT / AFTER_ROLLBACK / AFTER_COMPLETION四个阶段各自该放什么代码? - 什么时候进程内事件就够了、什么时候必须上 MQ?上了 MQ 之后「不丢、不重」分别靠什么保证?
先看一个几乎所有电商项目都写过的 OrderService。需求很朴素:「用户下单成功后,要发短信、加积分、发优惠券、记操作日志、通知仓库备货。」于是一位同学很自然地写了这样一段代码:
@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)); // 核心动作 —— 存订单 // 下面全是"顺带要做的事",却和核心动作挤在一起 smsService.send(cmd.userId(), "下单成功"); pointService.add(cmd.userId(), order.getAmount()); couponService.issue(cmd.userId(), order.getId()); auditLogService.record("PLACE_ORDER", order.getId()); warehouseService.notifyPrepare(order.getId()); return order; }}这段代码「能跑」,但它的构造器有 6 个依赖,方法体有 5 个与「创建订单」无关的调用。三个问题随之而来:
- 主流程被越拉越长:每加一个「下单后要干的事」,就要回来改
OrderService,它变成了一个什么都管的上帝类。 - 核心与非核心绑死:发优惠券的服务慢 3 秒,用户就要多等 3 秒;某个通知服务抛异常,整个下单直接回滚失败——非核心逻辑拖垮了核心业务。
- 事务边界混乱:短信在事务提交前就发出去了,一旦后面回滚,「下单成功的短信」永远收不回来。
把「下单」和「下单之后要做的事」解耦,正是事件驱动要解决的问题。核心业务只管发布一个事实(「订单已创建」),至于谁关心这个事实、要做什么,交给监听者各自决定。

这张流程图的左半(发布 → 监听器 → 事务阶段)是进程内的世界,右半(发到 broker → 幂等消费 → 补偿重试)是跨进程的世界。上半篇讲左半,下半篇讲右半;上面那张动图则是把左半的第一格放大到字节码级别。
Spring 的事件模型只有三个角色,记住它们就够用了:
| 角色 | 类型 / 注解 | 职责 |
|---|---|---|
| 事件 | 任意 POJO(Spring 4.2+ 不再强制继承 ApplicationEvent) | 描述「发生了什么」 |
| 发布者 | ApplicationEventPublisher | 把事件交给容器 |
| 监听者 | @EventListener 标注的方法 | 对事件做出反应 |
先定义事件。用 Java 17 的 record 最省事——它天然不可变,天然适合当「事实凭证」:
package com.example.order.event;import java.math.BigDecimal;import java.time.LocalDateTime;// 一个普通的不可变对象即可作为事件,无需继承任何基类public record OrderPlacedEvent(Long orderId, Long userId, BigDecimal amount, LocalDateTime placedAt) {}发布者只需注入 ApplicationEventPublisher,然后 publishEvent:
@Transactionalpublic Order placeOrder(PlaceOrderCommand cmd) { Order order = orderRepository.save(Order.create(cmd)); // 只表达"订单已创建"这个事实,不关心谁处理、怎么处理 publisher.publishEvent(new OrderPlacedEvent( order.getId(), cmd.userId(), order.getAmount(), LocalDateTime.now())); return order;}监听者用 @EventListener 接收,一个事件可以有任意多个监听器:
@Componentpublic class SmsListener { private static final Logger log = LoggerFactory.getLogger(SmsListener.class); @EventListener public void onOrderPlaced(OrderPlacedEvent event) { log.info("发送下单短信: userId={}, orderId={}", event.userId(), event.orderId()); // smsService.send(...) }}publishEvent只负责「广播」,发布者不知道有几个监听器- 监听器按方法参数类型匹配事件,参数是
OrderPlacedEvent就只收这个类型 - 新增一个「下单后要做的事」,只需新写一个监听器,
OrderService一行都不用改
说明:继承 ApplicationEvent 是 Spring 4.2 之前的老写法,那个时代事件必须继承基类、事件源要用构造器传入。现在任意对象都能当事件,record、普通 DTO 甚至 String 都可以,代码因此清爽许多。
先说结论:默认情况下,publishEvent 是在当前线程里同步执行的。听起来反直觉——明明是「发布/订阅」,怎么变成同步调用了?但这就是 Spring 的默认行为。用一个最小实验验证一下:
@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("发布事件,当前线程 = {}", Thread.currentThread().getName()); publisher.publishEvent(new OrderPlacedEvent(1L, 100L, new BigDecimal("9.9"), LocalDateTime.now())); log.info("publishEvent 返回,当前线程 = {}", 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("监听器执行,当前线程 = {}", Thread.currentThread().getName()); }}访问 /place,日志会是这样的顺序:
发布事件,当前线程 = http-nio-8080-exec-1监听器执行,当前线程 = http-nio-8080-exec-1publishEvent 返回,当前线程 = http-nio-8080-exec-1三行线程名完全相同,publishEvent 返回 排在监听器之后——证明监听器就在发布线程里同步跑完的。这带来两个必须记住的后果:
- 监听器耗时 = 请求耗时:监听器里做一次 3 秒的远程调用,用户的这次请求就要多等 3 秒。
- 异常会向上传播:某个监听器抛出异常,会直接冒泡到
publishEvent调用处,导致整个下单事务回滚。这正是「非核心逻辑拖垮核心业务」的真实机制。
很多人以为事件天然异步,于是把发短信、推送这类耗时操作直接塞进监听器,结果接口 RT 暴涨。请记住:不加 @Async 的监听器就是同步方法调用。
把这条结论用内核实验跑一遍——先看「同步发布到底阻塞了谁」:
再把监听器加到三个,看多播器的 for 循环与异常传播:
上面两格看的是「结果」。把同一段流程再摊成一次单步执行,一行一行走,才算真正看见「喊话到底发生在谁的线程上」——连点「下一步」,盯住线程名和累计耗时:
publishEvent(new OrderPlacedEvent(orderNo)); // ① 业务在喊话// ↓ ApplicationEventMulticaster 就地开播,不换线程smsListener.onOrderPlaced(e); // ② @Order(1) 短信先跑pointsListener.onOrderPlaced(e); // ③ @Order(2) 积分跟上auditListener.onOrderPlaced(e); // ④ @Order(3) 审计收尾// ⑤ 三行全部返回,publishEvent 才返回,一行抛出后面全跳过| 线程名 | http-nio-8080-exec-1 |
| 累计耗时 | 0 ms |
OrderController.payOrderService.placeOrder要让监听器脱离发布线程,需要 @EnableAsync 加上 @Async。但工程里不该裸用——@Async 默认用 SimpleAsyncTaskExecutor(每次新建线程,不复用),必须配一个正经线程池,并处理异步异常的「黑洞」问题:
@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-"); // 便于在日志里识别 executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); executor.initialize(); return executor; } // 异步方法抛出的异常不会传播到调用方,必须在这里兜底 @Override public AsyncUncaughtExceptionHandler getAsyncUncaughtExceptionHandler() { return (ex, method, params) -> log.error("异步事件执行失败: {}", method.getName(), ex); }}监听器上指定线程池即可异步执行:
@Componentpublic class PointListener { @Async("eventExecutor") @EventListener public void onOrderPlaced(OrderPlacedEvent event) { pointService.add(event.userId(), event.amount()); }}警告:异步监听器抛出的异常不会影响发布方的事务,这既是优点也是陷阱——出错时不会有任何反馈,如果不配 AsyncUncaughtExceptionHandler,异常会被静默吞掉,积分没加上你也不知道。
配上线程池再跑一次同样的事件,日志里的线程名会当场「搬家」:
异步的代价是异常没人接。这一分支专门演示异常去了哪里:
现在回到第一节那个事故:placeOrder 里发了短信,但事务随后回滚了,用户收到「下单成功」却是空欢喜。根因是事件的副作用发生在事务提交之前。@TransactionalEventListener 就是为此而生,它把监听器的执行时机精确绑定到事务的生命周期上:
| 阶段 | 触发时机 | 典型用途 |
|---|---|---|
BEFORE_COMMIT | 事务提交前 | 提交前校验、最后一道防线 |
AFTER_COMMIT | 事务成功提交后(默认) | 发短信、加积分、发消息——最常用 |
AFTER_ROLLBACK | 事务回滚后 | 记录失败、触发补偿、清缓存 |
AFTER_COMPLETION | 事务结束后(提交或回滚都会) | 无差别收尾,如释放资源 |

把四个阶段摊到一条时间线上看就清楚了:BEFORE_COMMIT 还在「可以反悔」的区域里,COMMIT 才是数据的生死线,而绝大多数副作用应该放在生死线之后。这条动画时间线也是排查「监听器到底跑没跑」的对照表——先确认你的发布点处于哪个阶段。

用代码对比一下「事故版」和「修复版」:
// 事故版:普通 @EventListener,在事务提交前执行@EventListenerpublic void onPlaced(OrderPlacedEvent event) { smsService.send(event.userId(), "下单成功"); // 事务若回滚,短信已发出,无法撤回}// 修复版:只有事务成功提交,才会执行@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)public void onCommitted(OrderPlacedEvent event) { smsService.send(event.userId(), "下单成功");}@TransactionalEventListener(phase = TransactionPhase.AFTER_ROLLBACK)public void onRolledBack(OrderPlacedEvent event) { log.warn("订单 {} 事务回滚,本次事件作废", event.orderId());}AFTER_COMMIT保证「数据落库了」才发短信,从根本上杜绝「短信先于数据」AFTER_ROLLBACK让你有机会对失败订单做补偿或告警- 同理,跨进程的 MQ 消息也应该在
AFTER_COMMIT里发送,否则同样会「消息先于数据」
坑:@TransactionalEventListener 有一个 fallbackExecution 开关,默认是 false。这意味着如果发布事件时当前根本没有事务(比如在 @Transactional 方法外部发布),监听器将完全不执行,而且不报错。排查时请先确认「发布点是否处于事务中」。确实需要在无事务时也执行,就显式写 @TransactionalEventListener(fallbackExecution = true)。
先亲手看一次「事件被暂存到提交之后才放行」的完整过程:
当同一个事件有多个监听器时,Spring 用 @Order 控制顺序(数字越小越先执行);用 SpEL 表达式做条件过滤;甚至可以让监听器返回一个对象继续发布新事件,形成事件链:
@Componentpublic class OrderEventListeners { // 只有金额大于 100 的订单才触发风控 @Order(1) @EventListener(condition = "#event.amount > 100") public void riskCheck(OrderPlacedEvent event) { riskService.check(event.orderId()); } // 返回非空对象 → Spring 会把它作为新事件继续发布(事件链) @Order(2) @EventListener public RewardEvent reward(OrderPlacedEvent event) { return new RewardEvent(event.userId(), event.amount()); }}@Order(1)让风控先于发奖执行condition = "#event.amount > 100"用 SpEL 引用事件参数,避免在方法体里写 if- 返回
RewardEvent会触发下一轮发布,RewardListener可以继续处理,形成链式反应
要点:顺序和条件只对同一线程内的同步监听器可靠。一旦用了 @Async,多个监听器并行执行,执行顺序不再有保证。
顺序这件事值得亲眼看一次:先按 @Order 跑一遍,再给其中一个监听器加上 @Async,你会发现「先审计后通知」的约定直接消失了。
另一个常见追问是:监听器里自己也要写库,该不该开一个新事务?把传播行为切成 REQUIRES_NEW 与 REQUIRED 各跑一次,看看主流程回滚时「已经写进去的那条记录」命运如何:
事件方案在单进程里非常好用,但分布式部署后立刻失效。假设你把服务部署了 3 个实例:
- 多实例收不到:A 实例发布的事件只存在于 A 的 JVM 内存里,B、C 实例的监听器永远收不到。
- 没有持久化:事件是内存对象,服务重启就丢失,没有重放、没有追溯。
- 没有重试与确认:监听器抛异常只有一次机会,没有投递确认、没有失败重试机制。
- 没有解耦跨服务:用户服务和积分服务分属不同应用时,进程内事件完全无能为力。
要让事件跨越进程边界,就需要引入消息队列(MQ):把「事件」变成一条持久化的消息,通过 broker 中转,各消费者进程各自订阅。

左边那列的每一项都是本节四条失效原因的对照:进程内事件快、零运维、能直接读到事务上下文,但出了这个 JVM 就什么都不是;MQ 换来持久化、重试、回放和跨服务能力,代价是网络往返、broker 运维,以及「投递语义只剩至少一次」——于是幂等变成了你的责任。
上面几格还是「别人按好的按钮」。想从命令行自己走一遍「喊话 → 直播 → 群消息」的完整转变,下面这台控制台连着同一个内核,回显全部由内核算出来:
把 lab event sync 和 lab event async 连着敲,对着看两次输出的线程名——一个是 http-nio-8080-exec-1,一个是 evt-1。名字一变,同步与异步的所有差别都从这一格长出来。
先加依赖,再配连接与可靠性参数。依赖不用背着写——勾一遍,看 pom 里什么该有、什么不该多:
<?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>依赖就位之后,连接与可靠性参数就是下面这几行——publisher-confirm-type 与手动 ACK 是「可靠投递」的前半场:
spring: rabbitmq: host: localhost port: 5672 username: guest password: guest publisher-confirm-type: correlated # 开启生产者确认 publisher-returns: true # 不可路由的消息回调 listener: simple: acknowledge-mode: manual # 手动 ACK,消费成功才确认 retry: enabled: true max-attempts: 3 # 本地重试 3 次RabbitMQ 的核心是「交换机按规则把消息投递到队列」,三种交换机类型要分清:
| 交换机类型 | 路由规则 | 典型场景 |
|---|---|---|
| Direct | 路由键精确匹配 | 点对点、按业务类型分发 |
| Topic | 路由键通配符匹配(* 一个词、# 多个词) | 事件总线,如 order.placed.# |
| Fanout | 忽略路由键,广播到所有绑定队列 | 广播通知、缓存刷新 |
用 Java Config 声明交换机、队列与绑定,再用 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) { // order.placed.xxx 的消息都会进入积分队列 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); }}用 @RabbitListener 消费,配合手动 ACK 保证「处理成功才确认」:
@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); // 处理成功,确认删除 } catch (Exception ex) { log.error("积分处理失败,将重新入队: orderId={}", event.orderId(), ex); channel.basicNack(deliveryTag, false, true); // 失败,重新入队 } }}注意:basicNack(deliveryTag, false, true) 会无限重入队,遇到「毒消息」会一直循环。生产上要么设置最大重试次数后转入死信队列,要么在确认前判断重试次数。
调消费者数量时先别急着手改配置——用沙盘把两个旋钮拨一拨,看积压和重复扣分分别长什么样:
queue depth: 1,200 ↓consume rate: 2,480 msg/sdedupe hits: 6 (skipped)lag ETA: 0.5s
顺带看一眼「一次慢请求在整条链路上到底卡在哪一格」——监听器就住在你的方法栈里,链路图上是找不到它的:
死信队列(DLX)像快递站那块「三次投递不成功就单独上架」的退货货架——包裹(消息)不再占用正常分拣通道,也不会凭空消失;你事后统一登记、查明原因、再决定重派还是销毁。没有这块货架时,站点只会被同一批退件反复堆满,新包裹彻底进不来。
两者都能解耦,但定位不同,选型要看业务特征:
| 维度 | RabbitMQ | Kafka |
|---|---|---|
| 吞吐 | 万级/秒 | 十万级~百万级/秒 |
| 延迟 | 低(毫秒级) | 低,但批量模型下略高 |
| 消息顺序 | 队列级有序 | 分区内有序 |
| 消息回放 | 消费即删,难以回放 | 保留期内可重复消费 |
| 模型 | 智能 broker + 简单队列 | 简单 broker + 智能消费者、日志式 |
| 适用场景 | 业务事件、任务分发、路由复杂 | 日志采集、流处理、大数据管道 |
一句话选型:业务事件(下单、支付这类「一件事触发几件副作用」)优先 RabbitMQ;海量日志、数据管道、需要回放或流式计算选 Kafka。
引入 MQ 后,新的难题叫「消息可靠性」——怎样保证消息不丢、不重。先看消息会在哪三个时机丢失:
| 丢失时机 | 场景 | 对策 |
|---|---|---|
| 生产前 | 业务事务已提交,但发消息前宕机 | 本地消息表:业务与消息同一事务写入 |
| 发送中 | 消息发出但 broker 未收到 | 生产者确认(publisher-confirm)+ 定时补偿 |
| 消费中 | 消费成功但 ACK 前宕机(重复投递) | 消费者幂等 + 手动 ACK |
核心思路是本地消息表:把「要发的消息」和业务数据写在同一个数据库事务里,保证两者要么都成功、要么都失败;再由一个定时任务扫描待发送的消息,可靠地投递到 MQ。
CREATE TABLE t_local_message ( id BIGINT PRIMARY KEY AUTO_INCREMENT, biz_type VARCHAR(64) NOT NULL COMMENT '业务类型,如 ORDER_PLACED', biz_id BIGINT NOT NULL COMMENT '业务主键,用于幂等', payload JSON NOT NULL COMMENT '消息体', status TINYINT NOT NULL DEFAULT 0 COMMENT '0待发送 1已发送 2已消费', retry_count INT NOT NULL DEFAULT 0 COMMENT '重试次数', next_retry_time DATETIME NOT NULL COMMENT '下次重试时间', create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, UNIQUE KEY uk_biz (biz_type, biz_id) -- 唯一键天然保证幂等);业务事务内同时写订单和消息,两者共用同一个事务:订单落库,消息就一定落库。
@Transactionalpublic Order placeOrder(PlaceOrderCommand cmd) { Order order = orderRepository.save(Order.create(cmd)); // 与订单在同一个事务里写入本地消息,保证"数据在则消息在" localMessageRepository.save(LocalMessage.pending( "ORDER_PLACED", order.getId(), toPayload(order))); return order;}定时任务扫描 status = 0 且到期未发的消息,投递成功后更新状态,失败则递增重试:
@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("消息补偿发送失败: id={}", message.getId(), ex); repository.bumpRetry(message.getId(), nextRetryTime()); } } }}发送端做到「至少一次」,消费端就必须做到幂等——用唯一键去重,重复消息直接跳过:
@Transactionalpublic void onOrderPlaced(OrderPlacedEvent event) { // 尝试插入幂等键,唯一键冲突说明这条消息已处理过 if (!dedupeRepository.tryInsert("ORDER_PLACED", event.orderId())) { log.info("检测到重复消息,跳过: orderId={}", event.orderId()); return; } pointService.add(event.userId(), event.amount());}要点:可靠投递 = 本地消息表(不丢)+ 生产者确认(发出有据)+ 定时补偿(失败重试)+ 消费者幂等(不重)。四件事凑齐,才能把「至少一次」的投递语义,转成对业务而言的「恰好一次」效果。
四件事不是四选一,而是一条流水线上的四个工位——缺哪个工位,就从哪个环节漏消息。把它们排进同一张图,比死记四个名词更有效:

左边两格守「发得出去」:本地消息表保证消息不会因为宕机而蒸发,生产者确认保证每条发出的消息都有回执;右边两格守「只生效一次」:定时补偿把失败的消息捞回来重试,消费者幂等让重复投递在业务侧归零。
@TransactionalEventListener 默认 fallbackExecution = false,只在存在事务时才会触发。如果你在事务外发布事件,监听器静默不执行,排查半天找不到原因。请先确认发布点被 @Transactional 包裹。
消费者没做幂等,是重复扣款的经典现场。MQ 的「至少一次」语义下,网络抖动会让同一条消息投递多次——没有幂等键,积分就会被加两次、库存就会被扣两次。幂等键必须是业务唯一键(如 orderId),不能是消息 id。
学到这里,你脑子里已经装了不少「正确说法」。但真正决定少踩坑的,是那些听起来都对、其实错的直觉。来玩一局辨伪——左边六句「直觉」,右边是它们的真相,配错会告诉你事故现场长什么样:
下面每一行的「现象原文」都可以整段复制去搜索。新手在这五个地方最容易卡住:监听器压根不跑、主流程被带崩、@Async 没生效、消息重复消费、消息堆积。
| 现象原文(片段) | 真实原因 | 30 秒自救 | 深挖看第几篇 |
|---|---|---|---|
@TransactionalEventListener 一行日志都不打,也不报错 | 发布点当前没有事务,而 fallbackExecution 默认是 false——容器直接把监听器跳过 | 确认调用链上有 @Transactional;确实要在无事务时执行就写 @TransactionalEventListener(fallbackExecution = true)。别用 @EventListener 试探,那会引入「提交前发短信」的新 bug | 本篇第五节 · 沙盘上方的 phase 实验 |
| 下单接口偶发 500,栈顶是监听器里的第三方调用,订单还整个回滚了 | 监听器默认同步执行且共用同一个事务,异常沿调用栈回抛到 publishEvent | 把非核心副作用改成 AFTER_COMMIT 阶段触发,并在方法体内 try/catch 自己吞掉;真要并行再配线程池 | 本篇第三、五节 · #17 @Transactional |
加了 @Async 但日志里线程名还是 http-nio-8080-exec-1 | 三种典型漏网:启动类缺 @EnableAsync;监听器和发布者在同一个 Bean 里(自调用不走代理);@Async 与 @EventListener 的方法签名是 void 之外还带了返回值 | 先看启动日志有没有 TaskExecutor 装配;把监听器挪到独立 @Component;异步监听器必须返回 void | #40 @Async 与线程池 |
org.springframework.core.task.SimpleAsyncTaskExecutor 出现在栈里,线程数一路涨到几千 | @Async 找不到候选线程池时的兜底实现——它每个任务 new 一个线程,不复用也不限流 | 显式 @Async("eventExecutor"),并注册一个 ThreadPoolTaskExecutor;Boot 里也可以配 spring.task.execution.pool.* | #40 @Async 与线程池 |
AsyncUncaughtExceptionHandler 从没被调用,异步监听器的异常凭空消失 | 该处理器只对返回 void 的异步方法生效;带 Future/返回值时异常被塞进结果对象,你不调 get() 就永远看不见 | 统一让事件监听器返回 void;或改用 CompletableFuture 并挂 .exceptionally(...) | 本篇第四节 · async(exc) 实验 |
用户积分被加了两次 / 库存被扣了两次,MQ 控制台显示同一条消息 redelivered=true | 「至少一次」语义下必然重投(ACK 超时、消费者重启、rebalance)。业务侧没做幂等 | 用业务唯一键(orderId)建幂等表 + 唯一索引,插入冲突即判定为重复并跳过;不要用消息 id 当幂等键 | 本篇第十节 · 上方沙盘 |
com.rabbitmq.client.ShutdownSignalException: channel error; reason: "NOT_FOUND - home node 'rabbit@...'" 或队列根本不存在 | @RabbitListener(queues = "...") 只声明消费方,不创建队列;交换机/队列/绑定三者少了一个 | 补 Queue / TopicExchange / Binding 三个 @Bean,或用 @QueueBinding 在注解上内联声明 | 本篇第八节 |
队列深度持续上涨(messages: 41932),消费者却显示 idle | 单条消息处理太慢(同步远程调用)、prefetch 太小、或全部并发消费者都卡在同一个慢 SQL 上 | 先 curl 看消费者线程栈;调大 concurrentConsumers 与 prefetch;把耗时操作拆成第二个队列异步处理 | 本篇第八节 · 消费端调参沙盘 |
大量消息进了死信队列 order.point.dlx,没人处理 | 重试次数用尽后按配置转 DLX,但你没写 DLX 消费者——毒消息从此沉睡 | 给 DLX 挂一个消费者做「落库 + 告警 + 人工重投」,并把最大重试次数与退避策略写进配置 | 本篇第八节注意条 |
上表里的报错都还能靠原文自救,但有一类事故不会这么善良:堆栈指向的行号看起来「没错」,凶手藏在更早的一行里。下面这份堆栈来自一次真实的「重复短信」事故——先别看结论,点出你认为的凶手帧:
上线第二天,客服转来三张截图:同一笔订单,客户连收三条「下单成功」短信。排查发现第一条处理的 ACK 在网络上丢了,broker 把消息重新投了两次。你从消费日志里捞到这第二次投递时抛出的堆栈——先别看结论,点出你认为的凶手帧。
目标:五分钟跑通「同步 → 异步 → AFTER_COMMIT」三种监听形态,用线程名和执行顺序证明它们真的不同。
第一步,pom.xml(只要 Web 与事务相关的最小集,不需要 MQ):
<?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>第二步,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"第三步,一张最小表 src/main/resources/schema.sql(用来制造「真的有事务」):
CREATE TABLE IF NOT EXISTS t_order ( id BIGINT PRIMARY KEY AUTO_INCREMENT, user_id BIGINT NOT NULL, amount DECIMAL(10,2) NOT NULL);第四步,事件、监听器与接口(包名 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("[同步] 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("[异步] 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={} 本次作废", 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.beans.factory.annotation.Value;import org.springframework.context.ApplicationEventPublisher;import org.springframework.jdbc.core.JdbcTemplate;import org.springframework.scheduling.annotation.EnableAsync;import org.springframework.scheduling.annotation.Async;import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;import org.springframework.context.annotation.Bean;import org.springframework.context.annotation.Configuration;import org.springframework.web.bind.annotation.*;import org.springframework.transaction.annotation.Transactional;@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("准备发布事件 thread={}", Thread.currentThread().getName()); publisher.publishEvent(new OrderPlacedEvent(id, 100L, new BigDecimal("9.90"))); if (boom) { throw new IllegalStateException("故意回滚"); } return Map.of("id", id, "status", "created"); }}第五步,跑起来并对比两条 curl:
# 正常下单:同步监听器与 AFTER_COMMIT 都会出现curl -s -X POST "http://localhost:8080/orders"# 故意回滚:同步监听器照样跑了(坑!),AFTER_COMMIT 没有,改走 AFTER_ROLLBACKcurl -s -X POST "http://localhost:8080/orders?boom=true"预期输出(第一次请求):
{"id":1,"status":"created"}14:32:05.112 [http-nio-8080-exec-1] INFO c.e.e.EventLabController - 准备发布事件 thread=http-nio-8080-exec-114:32:05.113 [http-nio-8080-exec-1] INFO c.e.e.LabListeners - [同步] orderId=1 thread=http-nio-8080-exec-114:32:05.114 [http-nio-8080-exec-1] INFO c.e.e.EventLabController - ... (控制器返回)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 - [异步] orderId=1 thread=evt-2 cost=509ms对着这张日志核对四件事:① [同步] 与发布者同线程且在返回之前;② [AFTER_COMMIT] 排在控制器之后(说明它晚于提交);③ [异步] 线程名变成 evt-2,时间戳最晚,控制器已经返回了它才结束;④ 第二次请求带 boom=true 时,[同步] 依然打印、[AFTER_COMMIT] 缺席、[AFTER_ROLLBACK] 出现——这就是第三节那个坑的铁证。
验收清单:能说出上面四条各自对应本文哪一节;能把 @TransactionalEventListener 的 phase 换成 BEFORE_COMMIT 并解释为什么此时仍能让短信发不出去。
只改一处,观察结论完全不同:
- 删掉控制器的
@Transactional,其余不动。你会观察到:[同步]照常打印,但[AFTER_COMMIT]再也不出现,也没有任何警告——亲手复现fallbackExecution = false的静默失效;再把 phase 注解改成@TransactionalEventListener(fallbackExecution = true),它回来了。 - 把
asyncListen上的@Async("eventExecutor")改成裸@Async。你会观察到:线程名前缀从evt-变成task-(Boot 的applicationTaskExecutor);若你把 executor Bean 整个删掉,就会看到SimpleAsyncTaskExecutor每次新建线程的原始行为。 - 在
syncListen里加Thread.sleep(3000)。你会观察到:curl的响应时间直接从几十毫秒涨到 3 秒以上,而boom=true那条会因为监听器耗时把事务连接占满——这就是「监听器耗时 = 接口耗时」的量化版。 - 把
@Order分别设为 1 和 2 加到两个同步监听器上,再给其中一个加@Async。你会观察到:加上异步后两者的相对次序不再稳定,重复请求十次能看到两种顺序交替出现。
提示:做完第 1 条再回到本篇第五节的 phase 实验,两边结论应当完全对得上。
给自己做一个「事件可观测性」小项目:一个 OrderPlacedEvent,三种投递方式(同步 / 异步 / AFTER_COMMIT + MQ 假实现),外加自动体检报告。
需求:
- 一个
ApplicationListener统计每类事件的发布次数与平均处理耗时(用System.nanoTime()包住publishEvent),每分钟打印一份 Top 5 - 一个
@RestController提供GET /events/replay?orderId=xxx,把某次事件重新投递给所有监听器,用于排障复现 - 一个「本地消息表」最小实现:
t_local_message表 +@Scheduled(fixedDelay = 5000)补偿任务 + 状态机0 待发送 → 1 已发送 → 2 已消费 - 消费端幂等:用
(biz_type, biz_id)唯一索引,捕获DuplicateKeyException视为已处理 - 不许修改业务监听器一行代码
验收清单:① 关掉 MQ(把 broker 地址写错)后压测 100 单,t_local_message 里最终 status=1 的记录数等于订单数,一条不少;② 手动重投同一 orderId 五次,积分只加一次且日志出现「检测到重复消息,跳过」;③ /events/replay 能复现出与原始日志一致的监听器顺序;④ 报告里的平均耗时与你手工测得的 P99 差距在 20% 以内。
不看上文,按顺序说出 publishEvent 之后多播器做了哪四件事(找监听器 → 排序 → 逐个调用 → 返回),并指出「异步」是在哪一步插进去的。
@EventListener、@Async @EventListener、@TransactionalEventListener(AFTER_COMMIT) 三种写法,各自的执行线程是什么?执行时机相对「事务提交」是早还是晚?
为什么说「幂等键必须是业务唯一键,不能用消息 id」?举一个具体事故场景说明。
本地消息表解决的是三段丢失时机中的哪一段?剩下两段分别由什么机制兜住?
什么信号一出现,就说明你应该停止使用进程内事件、改用真正的 MQ?至少说出三条。
喊话必同步,异步要池子;副作用等提交,跨楼才上路;上路必留痕,收件必幂等。
事件驱动的精髓是「解耦发生的事实与对它的反应」。进程内用 ApplicationEvent + @TransactionalEventListener(AFTER_COMMIT) 保证「数据先落库、副作用后执行」;跨进程用 MQ 把事件变成持久化消息,再以本地消息表 + 幂等消费守住可靠性。一句话记住:核心业务只发布事实,副作用各安其位,可靠投递靠"表 + 确认 + 补偿 + 幂等"四件套。