事件驱动:ApplicationEvent 与消息队列集成

bee2026-10-0893 分钟0 次阅读
从 Spring 内置事件讲到事务事件同步、@Async 异步事件,再到 RabbitMQ/Kafka 的取舍与可靠投递(本地消息表)方案。
1 / 163
小节
〇、30 秒看懂
2 / 163

先说人话:事件机制就是「把『发生了一件事实』和『对这件事做出反应』拆成两拨代码」。下单这个动作只负责宣布一条事实——「订单已创建」;至于是谁要为此发短信、谁要加积分、谁要通知仓库,下单代码一概不知道、也不关心。这样主流程就不会被一堆「顺带要做的事」越拖越长,也不会因为某个通知服务挂了而整单回滚。

3 / 163

有五个词后面会反复出现,先各给一句话解释:

4 / 163
  • 事件(Event):一个普通对象,只描述「发生了什么」,本身不带任何处理逻辑
  • 发布者(Publisher):调用 publishEvent(...) 把事件交给容器的那段业务代码
  • 监听器(Listener):标了 @EventListener 的方法,容器看到匹配的事件就调它
  • 多播器(Multicaster):容器内部那个「拿着名单挨个打电话」的对象,真正决定同步还是异步的是它
  • 消息队列(MQ):跨进程的邮局,把事件变成一条存下来的消息,由别的进程去取
5 / 163
类比

进程内事件像在办公室里喊一嗓子——你一开口,整间屋子的人都停下手里的事听完才继续干活(同步!你的方法栈里就卡着他们),而且隔壁楼的人根本听不见(跨不了进程)。消息队列像微信群通知——你发完就把手机扣下,谁有空谁看、没网的人上线还能翻到聊天记录(持久化),群消息还能一条条重放(可回放)。所以「要不要上 MQ」的本质问题是:这件事需要通知隔壁楼吗?

6 / 163
原理动画
动图 · 一次发布的同步全过程
动图 · 一次发布的同步全过程
7 / 163

上面这张动图是本篇最重要的一条主线:默认情况下 publishEvent 就是一根线程里的 for 循环。第七节的实验会让你亲眼看到三行日志的线程名一模一样。

8 / 163

学完这一篇,你应该能回答三个问题:

9 / 163
  • 为什么「加了 @EventListener」不等于「异步」?想让副作用离开主线程,到底缺了哪一块拼图?
  • 短信为什么必须在事务提交后才发?BEFORE_COMMIT / AFTER_COMMIT / AFTER_ROLLBACK / AFTER_COMPLETION 四个阶段各自该放什么代码?
  • 什么时候进程内事件就够了、什么时候必须上 MQ?上了 MQ 之后「不丢、不重」分别靠什么保证?
10 / 163
小节
一、为什么需要事件:从 5 个直连调用说起
11 / 163

先看一个几乎所有电商项目都写过的 OrderService。需求很朴素:「用户下单成功后,要发短信、加积分、发优惠券、记操作日志、通知仓库备货。」于是一位同学很自然地写了这样一段代码:

12 / 163
java
@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;    }}
13 / 163

这段代码「能跑」,但它的构造器有 6 个依赖,方法体有 5 个与「创建订单」无关的调用。三个问题随之而来:

14 / 163
  • 主流程被越拉越长:每加一个「下单后要干的事」,就要回来改 OrderService,它变成了一个什么都管的上帝类。
  • 核心与非核心绑死:发优惠券的服务慢 3 秒,用户就要多等 3 秒;某个通知服务抛异常,整个下单直接回滚失败——非核心逻辑拖垮了核心业务。
  • 事务边界混乱:短信在事务提交前就发出去了,一旦后面回滚,「下单成功的短信」永远收不回来。
15 / 163
提示

把「下单」和「下单之后要做的事」解耦,正是事件驱动要解决的问题。核心业务只管发布一个事实(「订单已创建」),至于谁关心这个事实、要做什么,交给监听者各自决定。

16 / 163
架构图
图 · 事件驱动的两种尺度
图 · 事件驱动的两种尺度
17 / 163

这张流程图的左半(发布 → 监听器 → 事务阶段)是进程内的世界,右半(发到 broker → 幂等消费 → 补偿重试)是跨进程的世界。上半篇讲左半,下半篇讲右半;上面那张动图则是把左半的第一格放大到字节码级别。

18 / 163
小节
二、Spring 事件三件套
19 / 163

Spring 的事件模型只有三个角色,记住它们就够用了:

20 / 163
对照表
角色类型 / 注解职责
事件任意 POJO(Spring 4.2+ 不再强制继承 ApplicationEvent)描述「发生了什么」
发布者ApplicationEventPublisher把事件交给容器
监听者@EventListener 标注的方法对事件做出反应
21 / 163

先定义事件。用 Java 17 的 record 最省事——它天然不可变,天然适合当「事实凭证」:

22 / 163
java
package com.example.order.event;import java.math.BigDecimal;import java.time.LocalDateTime;// 一个普通的不可变对象即可作为事件,无需继承任何基类public record OrderPlacedEvent(Long orderId, Long userId, BigDecimal amount, LocalDateTime placedAt) {}
23 / 163

发布者只需注入 ApplicationEventPublisher,然后 publishEvent:

24 / 163
java
@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;}
25 / 163

监听者用 @EventListener 接收,一个事件可以有任意多个监听器:

26 / 163
代码对照
代码java
@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 都可以,代码因此清爽许多。

27 / 163
小节
三、事件的同步本质:一个很多人不知道的坑
28 / 163

先说结论:默认情况下,publishEvent 是在当前线程里同步执行的。听起来反直觉——明明是「发布/订阅」,怎么变成同步调用了?但这就是 Spring 的默认行为。用一个最小实验验证一下:

29 / 163
java
@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";    }}
30 / 163
java
@Componentpublic class SmsListener {    private static final Logger log = LoggerFactory.getLogger(SmsListener.class);    @EventListener    public void onOrderPlaced(OrderPlacedEvent event) {        log.info("监听器执行,当前线程 = {}", Thread.currentThread().getName());    }}
31 / 163

访问 /place,日志会是这样的顺序:

32 / 163
text
发布事件,当前线程 = http-nio-8080-exec-1监听器执行,当前线程 = http-nio-8080-exec-1publishEvent 返回,当前线程 = http-nio-8080-exec-1
33 / 163

三行线程名完全相同,publishEvent 返回 排在监听器之后——证明监听器就在发布线程里同步跑完的。这带来两个必须记住的后果:

34 / 163
  • 监听器耗时 = 请求耗时:监听器里做一次 3 秒的远程调用,用户的这次请求就要多等 3 秒。
  • 异常会向上传播:某个监听器抛出异常,会直接冒泡到 publishEvent 调用处,导致整个下单事务回滚。这正是「非核心逻辑拖垮核心业务」的真实机制。
35 / 163
坑

很多人以为事件天然异步,于是把发短信、推送这类耗时操作直接塞进监听器,结果接口 RT 暴涨。请记住:不加 @Async 的监听器就是同步方法调用。

36 / 163

把这条结论用内核实验跑一遍——先看「同步发布到底阻塞了谁」:

37 / 163
内核实验
TeaVM办公室喊一嗓子:publishEvent 默认同步未启动
选「同步发布阻塞谁」,观察三行日志的线程名为什么完全相同
场景参数
点「运行演示」,在浏览器内真实执行 Java 编译出的内核算法,逐步看它怎么跑。
38 / 163

再把监听器加到三个,看多播器的 for 循环与异常传播:

39 / 163
内核实验
TeaVM多个监听器的顺序与「一人抛异常、全员停工」未启动
切到「多个监听器顺序」:@Order 小的先跑;某个监听器抛错时,后面的监听器和主流程一起失败
场景参数
点「运行演示」,在浏览器内真实执行 Java 编译出的内核算法,逐步看它怎么跑。
40 / 163

上面两格看的是「结果」。把同一段流程再摊成一次单步执行,一行一行走,才算真正看见「喊话到底发生在谁的线程上」——连点「下一步」,盯住线程名和累计耗时:

41 / 163
单步调试台
单步台单步走一遍 publishEvent:喊话到底发生在谁的线程上1 / 6
连点「下一步」六次:前三步共用同一个线程名,然后看监听器如何排队与成败
被调试的代码
1publishEvent(new OrderPlacedEvent(orderNo)); // ① 业务在喊话
2// ↓ ApplicationEventMulticaster 就地开播,不换线程
3smsListener.onOrderPlaced(e); // ② @Order(1) 短信先跑
4pointsListener.onOrderPlaced(e); // ③ @Order(2) 积分跟上
5auditListener.onOrderPlaced(e); // ④ @Order(3) 审计收尾
6// ⑤ 三行全部返回,publishEvent 才返回,一行抛出后面全跳过
此刻的变量
线程名http-nio-8080-exec-1
累计耗时0 ms
调用栈
1OrderController.pay
2OrderService.placeOrder
1这一行看着像「发出去」,其实是「当场喊」:事件只是被交给了多播器,还没有任何监听器开工,而且你还站在自己的线程上。
42 / 163
小节
四、@Async 异步事件:让副作用离开主线程
43 / 163

要让监听器脱离发布线程,需要 @EnableAsync 加上 @Async。但工程里不该裸用——@Async 默认用 SimpleAsyncTaskExecutor(每次新建线程,不复用),必须配一个正经线程池,并处理异步异常的「黑洞」问题:

44 / 163
java
@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);    }}
45 / 163

监听器上指定线程池即可异步执行:

46 / 163
代码对照
代码java
@Componentpublic class PointListener {    @Async("eventExecutor")    @EventListener    public void onOrderPlaced(OrderPlacedEvent event) {        pointService.add(event.userId(), event.amount());    }}
解读

警告:异步监听器抛出的异常不会影响发布方的事务,这既是优点也是陷阱——出错时不会有任何反馈,如果不配 AsyncUncaughtExceptionHandler,异常会被静默吞掉,积分没加上你也不知道。

47 / 163

配上线程池再跑一次同样的事件,日志里的线程名会当场「搬家」:

48 / 163
内核实验
TeaVM从「喊一嗓子」到「群消息」:@Async 监听器未启动
切到「@Async 监听器」:对比监听器线程 evt-1 与发布线程 http-nio-8080-exec-1
场景参数
点「运行演示」,在浏览器内真实执行 Java 编译出的内核算法,逐步看它怎么跑。
49 / 163

异步的代价是异常没人接。这一分支专门演示异常去了哪里:

50 / 163
内核实验
TeaVM异步异常的黑洞:为什么积分没加却一声不响未启动
选「异常去了哪里」:没有 AsyncUncaughtExceptionHandler 时,异常只落在日志里,调用方毫无感知
场景参数
点「运行演示」,在浏览器内真实执行 Java 编译出的内核算法,逐步看它怎么跑。
51 / 163
小节
五、@TransactionalEventListener:事务提交后再触发(重头戏)
52 / 163

现在回到第一节那个事故:placeOrder 里发了短信,但事务随后回滚了,用户收到「下单成功」却是空欢喜。根因是事件的副作用发生在事务提交之前。@TransactionalEventListener 就是为此而生,它把监听器的执行时机精确绑定到事务的生命周期上:

53 / 163
对照表
阶段触发时机典型用途
BEFORE_COMMIT事务提交前提交前校验、最后一道防线
AFTER_COMMIT事务成功提交后(默认)发短信、加积分、发消息——最常用
AFTER_ROLLBACK事务回滚后记录失败、触发补偿、清缓存
AFTER_COMPLETION事务结束后(提交或回滚都会)无差别收尾,如释放资源
54 / 163
原理动画
动图 · 事务事件的四个阶段时间线
动图 · 事务事件的四个阶段时间线
55 / 163

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

56 / 163
原理动画
动图 · 事务事件的四个阶段
动图 · 事务事件的四个阶段
57 / 163

用代码对比一下「事故版」和「修复版」:

58 / 163
代码对照
代码java
// 事故版:普通 @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)。

59 / 163

先亲手看一次「事件被暂存到提交之后才放行」的完整过程:

60 / 163
内核实验
TeaVMAFTER_COMMIT 到底几点触发未启动
切到「AFTER_COMMIT 时机」:对照 BEFORE_COMMIT / AFTER_COMMIT / AFTER_ROLLBACK / AFTER_COMPLETION 四个落点
场景参数
点「运行演示」,在浏览器内真实执行 Java 编译出的内核算法,逐步看它怎么跑。
61 / 163
内核实验
TeaVM事务边界决定事件何时可见未启动
用「外层回滚」演示:事务没提交,事件与数据会怎样
场景参数
点「运行演示」,在浏览器内真实执行 Java 编译出的内核算法,逐步看它怎么跑。
62 / 163
小节
六、监听器顺序与条件:精细化编排
63 / 163

当同一个事件有多个监听器时,Spring 用 @Order 控制顺序(数字越小越先执行);用 SpEL 表达式做条件过滤;甚至可以让监听器返回一个对象继续发布新事件,形成事件链:

64 / 163
代码对照
代码java
@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,多个监听器并行执行,执行顺序不再有保证。

65 / 163

顺序这件事值得亲眼看一次:先按 @Order 跑一遍,再给其中一个监听器加上 @Async,你会发现「先审计后通知」的约定直接消失了。

66 / 163
内核实验
TeaVM编排实验:@Order、condition 与异步带来的顺序崩塌未启动
先看「多个监听器顺序」的稳定次序,再切「@Async 监听器」比较输出顺序的变化
场景参数
点「运行演示」,在浏览器内真实执行 Java 编译出的内核算法,逐步看它怎么跑。
67 / 163

另一个常见追问是:监听器里自己也要写库,该不该开一个新事务?把传播行为切成 REQUIRES_NEW 与 REQUIRED 各跑一次,看看主流程回滚时「已经写进去的那条记录」命运如何:

68 / 163
内核实验
TeaVM监听器里的事务归属:REQUIRED vs REQUIRES_NEW未启动
选 REQUIRES_NEW:外层回滚后监听器的写入仍然保留——审计日志常这么配
场景参数
点「运行演示」,在浏览器内真实执行 Java 编译出的内核算法,逐步看它怎么跑。
69 / 163
小节
七、从单体事件到消息队列:进程内事件的边界
70 / 163

事件方案在单进程里非常好用,但分布式部署后立刻失效。假设你把服务部署了 3 个实例:

71 / 163
  • 多实例收不到:A 实例发布的事件只存在于 A 的 JVM 内存里,B、C 实例的监听器永远收不到。
  • 没有持久化:事件是内存对象,服务重启就丢失,没有重放、没有追溯。
  • 没有重试与确认:监听器抛异常只有一次机会,没有投递确认、没有失败重试机制。
  • 没有解耦跨服务:用户服务和积分服务分属不同应用时,进程内事件完全无能为力。
72 / 163

要让事件跨越进程边界,就需要引入消息队列(MQ):把「事件」变成一条持久化的消息,通过 broker 中转,各消费者进程各自订阅。

73 / 163
架构图
图 · 进程内事件 vs 消息队列
图 · 进程内事件 vs 消息队列
74 / 163

左边那列的每一项都是本节四条失效原因的对照:进程内事件快、零运维、能直接读到事务上下文,但出了这个 JVM 就什么都不是;MQ 换来持久化、重试、回放和跨服务能力,代价是网络往返、broker 运维,以及「投递语义只剩至少一次」——于是幂等变成了你的责任。

75 / 163
内核实验
TeaVM同一件事,两种尺度:喊一嗓子 vs 群通知未启动
切到「换成消息队列」:观察事件如何被序列化成消息、由 broker 投给另一个进程
场景参数
点「运行演示」,在浏览器内真实执行 Java 编译出的内核算法,逐步看它怎么跑。
76 / 163

上面几格还是「别人按好的按钮」。想从命令行自己走一遍「喊话 → 直播 → 群消息」的完整转变,下面这台控制台连着同一个内核,回显全部由内核算出来:

77 / 163
内核控制台
78 / 163
提示

把 lab event sync 和 lab event async 连着敲,对着看两次输出的线程名——一个是 http-nio-8080-exec-1,一个是 evt-1。名字一变,同步与异步的所有差别都从这一格长出来。

79 / 163
小节
八、RabbitMQ 整合速通
80 / 163

先加依赖,再配连接与可靠性参数。依赖不用背着写——勾一遍,看 pom 里什么该有、什么不该多:

81 / 163
生成器
生成器RabbitMQ 整合到底要哪几个 starterpom.xml2 / 3
只勾 web 得到空壳;补上 amqp 才有 RabbitTemplate 与 @RabbitListener 的自动配置;再带上 actuator,broker 的连通状态会直接摊到 /actuator/health 里——对着第七节的边界与本文后面的可靠投递读
产物
<?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>
勾了这些,代价与理由在这里
parent继承 3.3.4 的 starter-parent 之后,所有 spring-boot-starter-* 都不用写版本号;一旦有人手写给某个 starter 加 version,就以那条为准——这是依赖版本漂移最常见的原因。
Web做接口就绕不开它: DispatcherServlet、内嵌 Tomcat、JSON 序列化全在这个 starter 里。
AMQPRabbitTemplate 与监听容器;序列化消息要用 SpringMessagingConverter。
82 / 163

依赖就位之后,连接与可靠性参数就是下面这几行——publisher-confirm-type 与手动 ACK 是「可靠投递」的前半场:

83 / 163
yaml
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 次
84 / 163

RabbitMQ 的核心是「交换机按规则把消息投递到队列」,三种交换机类型要分清:

85 / 163
对照表
交换机类型路由规则典型场景
Direct路由键精确匹配点对点、按业务类型分发
Topic路由键通配符匹配(* 一个词、# 多个词)事件总线,如 order.placed.#
Fanout忽略路由键,广播到所有绑定队列广播通知、缓存刷新
86 / 163

用 Java Config 声明交换机、队列与绑定,再用 RabbitTemplate 发送:

87 / 163
java
@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.#");    }}
88 / 163
java
@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);    }}
89 / 163

用 @RabbitListener 消费,配合手动 ACK 保证「处理成功才确认」:

90 / 163
代码对照
代码java
@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) 会无限重入队,遇到「毒消息」会一直循环。生产上要么设置最大重试次数后转入死信队列,要么在确认前判断重试次数。

91 / 163

调消费者数量时先别急着手改配置——用沙盘把两个旋钮拨一拨,看积压和重复扣分分别长什么样:

92 / 163
沙盘
沙盘消费端调参沙盘:并发数 × 幂等开关
运行结果
queue depth: 1,200 ↓
consume rate: 2,480 msg/s
dedupe hits: 6 (skipped)
lag ETA: 0.5s
这是目标形态:并发够、幂等在,重复投递被静默跳过
93 / 163

顺带看一眼「一次慢请求在整条链路上到底卡在哪一格」——监听器就住在你的方法栈里,链路图上是找不到它的:

94 / 163
内核实验
TeaVM接口 RT 涨了 800ms,凶手是监听器未启动
选「慢请求与超时」:注意耗时统计里根本没有 multicastEvent 这一格,因为它就在业务方法内部
场景参数
点「运行演示」,在浏览器内真实执行 Java 编译出的内核算法,逐步看它怎么跑。
95 / 163
类比

死信队列(DLX)像快递站那块「三次投递不成功就单独上架」的退货货架——包裹(消息)不再占用正常分拣通道,也不会凭空消失;你事后统一登记、查明原因、再决定重派还是销毁。没有这块货架时,站点只会被同一批退件反复堆满,新包裹彻底进不来。

96 / 163
小节
九、Kafka 与 RabbitMQ 的取舍
97 / 163

两者都能解耦,但定位不同,选型要看业务特征:

98 / 163
对照表
维度RabbitMQKafka
吞吐万级/秒十万级~百万级/秒
延迟低(毫秒级)低,但批量模型下略高
消息顺序队列级有序分区内有序
消息回放消费即删,难以回放保留期内可重复消费
模型智能 broker + 简单队列简单 broker + 智能消费者、日志式
适用场景业务事件、任务分发、路由复杂日志采集、流处理、大数据管道
99 / 163

一句话选型:业务事件(下单、支付这类「一件事触发几件副作用」)优先 RabbitMQ;海量日志、数据管道、需要回放或流式计算选 Kafka。

100 / 163
小节
十、可靠投递:本地消息表方案(重点)
101 / 163

引入 MQ 后,新的难题叫「消息可靠性」——怎样保证消息不丢、不重。先看消息会在哪三个时机丢失:

102 / 163
对照表
丢失时机场景对策
生产前业务事务已提交,但发消息前宕机本地消息表:业务与消息同一事务写入
发送中消息发出但 broker 未收到生产者确认(publisher-confirm)+ 定时补偿
消费中消费成功但 ACK 前宕机(重复投递)消费者幂等 + 手动 ACK
103 / 163

核心思路是本地消息表:把「要发的消息」和业务数据写在同一个数据库事务里,保证两者要么都成功、要么都失败;再由一个定时任务扫描待发送的消息,可靠地投递到 MQ。

104 / 163
sql
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)   -- 唯一键天然保证幂等);
105 / 163

业务事务内同时写订单和消息,两者共用同一个事务:订单落库,消息就一定落库。

106 / 163
java
@Transactionalpublic Order placeOrder(PlaceOrderCommand cmd) {    Order order = orderRepository.save(Order.create(cmd));    // 与订单在同一个事务里写入本地消息,保证"数据在则消息在"    localMessageRepository.save(LocalMessage.pending(            "ORDER_PLACED", order.getId(), toPayload(order)));    return order;}
107 / 163

定时任务扫描 status = 0 且到期未发的消息,投递成功后更新状态,失败则递增重试:

108 / 163
java
@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());            }        }    }}
109 / 163

发送端做到「至少一次」,消费端就必须做到幂等——用唯一键去重,重复消息直接跳过:

110 / 163
代码对照
代码java
@Transactionalpublic void onOrderPlaced(OrderPlacedEvent event) {    // 尝试插入幂等键,唯一键冲突说明这条消息已处理过    if (!dedupeRepository.tryInsert("ORDER_PLACED", event.orderId())) {        log.info("检测到重复消息,跳过: orderId={}", event.orderId());        return;    }    pointService.add(event.userId(), event.amount());}
解读

要点:可靠投递 = 本地消息表(不丢)+ 生产者确认(发出有据)+ 定时补偿(失败重试)+ 消费者幂等(不重)。四件事凑齐,才能把「至少一次」的投递语义,转成对业务而言的「恰好一次」效果。

111 / 163

四件事不是四选一,而是一条流水线上的四个工位——缺哪个工位,就从哪个环节漏消息。把它们排进同一张图,比死记四个名词更有效:

112 / 163
架构图
图 · 可靠投递四件套
图 · 可靠投递四件套
113 / 163

左边两格守「发得出去」:本地消息表保证消息不会因为宕机而蒸发,生产者确认保证每条发出的消息都有回执;右边两格守「只生效一次」:定时补偿把失败的消息捞回来重试,消费者幂等让重复投递在业务侧归零。

114 / 163
小节
十一、两个高频坑 + 决策 + 总结
115 / 163
坑

@TransactionalEventListener 默认 fallbackExecution = false,只在存在事务时才会触发。如果你在事务外发布事件,监听器静默不执行,排查半天找不到原因。请先确认发布点被 @Transactional 包裹。

116 / 163
坑

消费者没做幂等,是重复扣款的经典现场。MQ 的「至少一次」语义下,网络抖动会让同一条消息投递多次——没有幂等键,积分就会被加两次、库存就会被扣两次。幂等键必须是业务唯一键(如 orderId),不能是消息 id。

117 / 163
决策
决策下单成功后要发短信、加积分,你对「订单已创建」这个事件的副作用应该用进程内事件还是 MQ?
118 / 163

学到这里,你脑子里已经装了不少「正确说法」。但真正决定少踩坑的,是那些听起来都对、其实错的直觉。来玩一局辨伪——左边六句「直觉」,右边是它们的真相,配错会告诉你事故现场长什么样:

119 / 163
配对闯关
闯关六个直觉说法,配它的真相已配对 0/6 · 配错 0
两列都打乱了;判据是「这句话在哪个环节会翻车」,配错时它会给你对应的事故现场
先点左边一个
120 / 163
小节
十二、常见报错速查
121 / 163

下面每一行的「现象原文」都可以整段复制去搜索。新手在这五个地方最容易卡住:监听器压根不跑、主流程被带崩、@Async 没生效、消息重复消费、消息堆积。

122 / 163
对照表
现象原文(片段)真实原因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 挂一个消费者做「落库 + 告警 + 人工重投」,并把最大重试次数与退避策略写进配置本篇第八节注意条
123 / 163

上表里的报错都还能靠原文自救,但有一类事故不会这么善良:堆栈指向的行号看起来「没错」,凶手藏在更早的一行里。下面这份堆栈来自一次真实的「重复短信」事故——先别看结论,点出你认为的凶手帧:

124 / 163
报错急救
报错急救DuplicateKeyException:幂等台账自己撞了唯一键
同一条消息投了三次,短信发了三条:凶手在监听器的第 24 行

上线第二天,客服转来三张截图:同一笔订单,客户连收三条「下单成功」短信。排查发现第一条处理的 ACK 在网络上丢了,broker 把消息重新投了两次。你从消费日志里捞到这第二次投递时抛出的堆栈——先别看结论,点出你认为的凶手帧。

org.springframework.dao.DuplicateKeyException: PreparedStatementCallback; SQL [INSERT INTO t_msg_dedupe(biz_type, biz_id) VALUES (?, ?)]; Duplicate entry 'ORDER_PLACED-88017' for key 'uk_biz'
at org.springframework.jdbc.support.SQLErrorCodeSQLExceptionTranslator.doTranslate(SQLErrorCodeSQLExceptionTranslator.java:246)
at org.springframework.jdbc.core.JdbcTemplate.update(JdbcTemplate.java:794)
at com.bee.pay.SmsListener.onOrderPlaced(SmsListener.java:24)
at org.springframework.amqp.rabbit.listener.adapter.MessagingMessageListenerAdapter.invokeHandler(MessagingMessageListenerAdapter.java:281)
at org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer.invokeListener(SimpleMessageListenerContainer.java:1473)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:659)
点你认为的「凶手行」(可反复试)
不会也没关系:先猜异常名,再猜哪一行在做决定。
125 / 163
随堂自测
随堂自测你在 placeOrder(标了 @Transactional)里 publishEvent 一个 OrderPlacedEvent,监听器写成 @TransactionalEventListener。压测时发现监听器从来没执行过,也没有任何报错。最可能的原因是?
先自己选一个,选中立刻告诉你对不对
126 / 163
随堂自测
随堂自测上线后用户反馈「收到下单成功短信,但订单查询里没有这一单」。排查发现监听器用的是普通 @EventListener 直接发短信。下列修法中最稳妥的是?
先自己选一个,选中立刻告诉你对不对
127 / 163
小节
十三、动手练习
128 / 163
小节
第一档 · 照做
129 / 163

目标:五分钟跑通「同步 → 异步 → AFTER_COMMIT」三种监听形态,用线程名和执行顺序证明它们真的不同。

130 / 163

第一步,pom.xml(只要 Web 与事务相关的最小集,不需要 MQ):

131 / 163
xml
<?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>
132 / 163

第二步,src/main/resources/application.yml:

133 / 163
yaml
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"
134 / 163

第三步,一张最小表 src/main/resources/schema.sql(用来制造「真的有事务」):

135 / 163
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);
136 / 163

第四步,事件、监听器与接口(包名 com.example.eventlab):

137 / 163
java
package com.example.eventlab;import java.math.BigDecimal;public record OrderPlacedEvent(Long orderId, Long userId, BigDecimal amount) {}
138 / 163
java
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());    }}
139 / 163
java
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");    }}
140 / 163

第五步,跑起来并对比两条 curl:

141 / 163
bash
# 正常下单:同步监听器与 AFTER_COMMIT 都会出现curl -s -X POST "http://localhost:8080/orders"# 故意回滚:同步监听器照样跑了(坑!),AFTER_COMMIT 没有,改走 AFTER_ROLLBACKcurl -s -X POST "http://localhost:8080/orders?boom=true"
142 / 163

预期输出(第一次请求):

143 / 163
json
{"id":1,"status":"created"}
144 / 163
text
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
145 / 163

对着这张日志核对四件事:① [同步] 与发布者同线程且在返回之前;② [AFTER_COMMIT] 排在控制器之后(说明它晚于提交);③ [异步] 线程名变成 evt-2,时间戳最晚,控制器已经返回了它才结束;④ 第二次请求带 boom=true 时,[同步] 依然打印、[AFTER_COMMIT] 缺席、[AFTER_ROLLBACK] 出现——这就是第三节那个坑的铁证。

146 / 163

验收清单:能说出上面四条各自对应本文哪一节;能把 @TransactionalEventListener 的 phase 换成 BEFORE_COMMIT 并解释为什么此时仍能让短信发不出去。

147 / 163
小节
第二档 · 变体
148 / 163

只改一处,观察结论完全不同:

149 / 163
  1. 删掉控制器的 @Transactional,其余不动。你会观察到:[同步] 照常打印,但 [AFTER_COMMIT] 再也不出现,也没有任何警告——亲手复现 fallbackExecution = false 的静默失效;再把 phase 注解改成 @TransactionalEventListener(fallbackExecution = true),它回来了。
  2. 把 asyncListen 上的 @Async("eventExecutor") 改成裸 @Async。你会观察到:线程名前缀从 evt- 变成 task-(Boot 的 applicationTaskExecutor);若你把 executor Bean 整个删掉,就会看到 SimpleAsyncTaskExecutor 每次新建线程的原始行为。
  3. 在 syncListen 里加 Thread.sleep(3000)。你会观察到:curl 的响应时间直接从几十毫秒涨到 3 秒以上,而 boom=true 那条会因为监听器耗时把事务连接占满——这就是「监听器耗时 = 接口耗时」的量化版。
  4. 把 @Order 分别设为 1 和 2 加到两个同步监听器上,再给其中一个加 @Async。你会观察到:加上异步后两者的相对次序不再稳定,重复请求十次能看到两种顺序交替出现。
150 / 163

提示:做完第 1 条再回到本篇第五节的 phase 实验,两边结论应当完全对得上。

151 / 163
小节
第三档 · 造一个
152 / 163

给自己做一个「事件可观测性」小项目:一个 OrderPlacedEvent,三种投递方式(同步 / 异步 / AFTER_COMMIT + MQ 假实现),外加自动体检报告。

153 / 163

需求:

154 / 163
  • 一个 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 视为已处理
  • 不许修改业务监听器一行代码
155 / 163

验收清单:① 关掉 MQ(把 broker 地址写错)后压测 100 单,t_local_message 里最终 status=1 的记录数等于订单数,一条不少;② 手动重投同一 orderId 五次,积分只加一次且日志出现「检测到重复消息,跳过」;③ /events/replay 能复现出与原始日志一致的监听器顺序;④ 报告里的平均耗时与你手工测得的 P99 差距在 20% 以内。

156 / 163
小节
十四、要点自查
157 / 163
自检

不看上文,按顺序说出 publishEvent 之后多播器做了哪四件事(找监听器 → 排序 → 逐个调用 → 返回),并指出「异步」是在哪一步插进去的。

158 / 163
自检

@EventListener、@Async @EventListener、@TransactionalEventListener(AFTER_COMMIT) 三种写法,各自的执行线程是什么?执行时机相对「事务提交」是早还是晚?

159 / 163
自检

为什么说「幂等键必须是业务唯一键,不能用消息 id」?举一个具体事故场景说明。

160 / 163
自检

本地消息表解决的是三段丢失时机中的哪一段?剩下两段分别由什么机制兜住?

161 / 163
自检

什么信号一出现,就说明你应该停止使用进程内事件、改用真正的 MQ?至少说出三条。

162 / 163
口诀

喊话必同步,异步要池子;副作用等提交,跨楼才上路;上路必留痕,收件必幂等。

163 / 163
总结

事件驱动的精髓是「解耦发生的事实与对它的反应」。进程内用 ApplicationEvent + @TransactionalEventListener(AFTER_COMMIT) 保证「数据先落库、副作用后执行」;跨进程用 MQ 把事件变成持久化消息,再以本地消息表 + 幂等消费守住可靠性。一句话记住:核心业务只发布事实,副作用各安其位,可靠投递靠"表 + 确认 + 补偿 + 幂等"四件套。