阶段 6 · Spring Boot 3.2.x / JDK 17
提升系统吞吐、解耦与可靠性的关键。生产环境里,MQ
和线程池用错是高频事故源。
1. 导语
一个下单接口,为什么在流量上来后响应变慢?因为你在一个请求里同步做了太多事:扣库存、发短信、发邮件、写操作日志、更新用户积分……其中任何一个慢操作(比如第三方短信服务偶发超时),都会拖慢整个接口的响应时间,甚至拖垮整个服务的线程池。
消息队列(Message
Queue,MQ)和异步任务(@Async)正是解决这类问题的两把利器:把非核心的耗时操作从主流程里剥离出去,让主流程快速返回,让后台慢慢处理。MQ
还能把「订单服务」和「库存服务」「通知服务」解耦——订单服务只需要发一条消息,至于谁消费、什么时候消费、消费失败怎么处理,订单服务一概不关心。这就是「削峰填谷」:双十一零点涌入的
10 万下单请求,不会瞬间打到数据库,而是先进入
MQ,由消费端按自己的节奏慢慢处理。
本篇是系列教程的第 6 阶段。在前面的阶段里,你已经掌握了 Web
层(Controller/Service/Mapper 分层)、数据访问(MyBatis-Plus /
JPA)、缓存(Redis)等能力。本篇要补上的是「系统异步化」这关键一环,并且把「消息可靠性」讲透——这是面试和生产的双重高频考点。
学完本篇,你能独立做到:
- 用 RabbitMQ 实现「下单 → 发消息 →
库存/通知服务异步消费」的完整链路; - 用死信队列(DLX)+ TTL 实现「订单 30 分钟未支付自动关闭」;
- 用 Kafka 实现高吞吐消费,并做到手动提交 offset + 幂等消费;
- 用
@Async+
自定义线程池优雅地异步发邮件,并能解释它为什么有时会「失效」; - 用「本地消息表」保证跨服务消息最终一致。
MQ
和线程池在生产环境是高频事故源——消息丢了、消息重复消费、@Async
不生效、定时任务多实例重复跑、线程池队列堆积导致
OOM……这些问题,本篇都会讲透原理,并给出可直接上生产的代码。
2. 学习目标与前置要求
学完你能…
- 说出消息队列的「三大可靠性问题」(丢失 / 重复 /
顺序)分别如何解决,并能画出解决方案对照表; - 独立写出 RabbitMQ 的生产者确认(confirm)、手动
ACK、死信队列延迟任务三段完整代码; - 独立写出 Kafka 手动提交 offset + 幂等消费的完整代码,并说清
Partition / Consumer Group / Offset 三者关系; - 独立配置一个自定义线程池并正确使用
@Async,能验证并解释「同类自调用导致失效」的原因; - 独立写出一个带 Redis 分布式锁的定时任务,避免多实例重复执行;
- 独立实现「本地消息表」模式,保证业务数据和消息数据同库同事务、消息不丢。
前置依赖
- 阶段 4(数据访问):理解 MyBatis-Plus / JPA 的
Entity、Mapper
用法,本篇实战项目会大量用到;若未掌握,建议先看《数据访问》篇。 - 阶段 4(数据访问):理解 Redis 基本命令与
StringRedisTemplate,本篇「幂等去重」「分布式锁」依赖它;若未掌握,建议先看《数据访问》篇。 - 阶段 1(Spring 核心与 AOP):理解 Spring AOP
代理机制(JDK 动态代理 / CGLIB),这是理解
@Async、@Transactional
失效的钥匙;若未掌握,建议先看《Spring 核心》篇。
3.
环境准备(本篇需要 RabbitMQ、Kafka、Redis、MySQL)
本篇代码依赖四个中间件,建议用 Docker
一键拉起(版本已固定,避免环境差异):
# 1. RabbitMQ 3.12(带管理界面,管理端口 15672)
docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672
rabbitmq:3.12-management
# 2. Kafka 3.6(KRaft 单节点模式,无需 ZooKeeper)
docker run -d --name kafka -p 9092:9092
-e KAFKA_CFG_NODE_ID=0
-e KAFKA_CFG_PROCESS_ROLES=controller,broker
-e KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=0@kafka:9093
-e KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093
-e KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092
-e KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER
-e KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
bitnami/kafka:3.6
# 3. Redis 7
docker run -d --name redis -p 6379:6379 redis:7
# 4. MySQL 8.0(密码 root / root,库名 demo)
docker run -d --name mysql -p 3306:3306
-e MYSQL_ROOT_PASSWORD=root -e MYSQL_DATABASE=demo mysql:8.0
项目依赖(pom.xml 关键片段,Spring Boot
3.2.x):
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>3.2.5</version>
</parent>
<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-amqp</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<dependency>
<groupId>com.baomidou</groupId>
<artifactId>mybatis-plus-spring-boot3-starter</artifactId>
<version>3.5.7</version>
</dependency>
<dependency>
<groupId>com.mysql</groupId>
<artifactId>mysql-connector-j</artifactId>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
<!-- fastjson2:消息体 JSON 序列化(本地消息表 payload、事件对象) -->
<dependency>
<groupId>com.alibaba.fastjson2</groupId>
<artifactId>fastjson2</artifactId>
<version>2.0.51</version>
</dependency>
</dependencies>
说明:
ApiResult、ErrorCode、BizException、GlobalExceptionHandler
等公共类是贯穿全系列的「统一基础设施」,定义与前面阶段逐字一致(见阶段
1)。本篇「生产级实战项目」章节末尾会附上完整定义,方便你独立运行;这些类在各篇文章里保持不变,请勿改名或重定义。
4. 正文章节
第 1 章 MQ 基础与选型
1.1 为什么需要消息队列
先回答一个最朴素的问题:没有
MQ,我们的系统难道不能跑吗?当然能。一个单体应用里,下单接口同步完成「校验
→ 扣库存 → 发短信 → 发邮件 →
写日志」,逻辑清晰、调试方便。但随着流量增长,这套同步写法的三个致命问题会逐渐暴露:
问题一:慢操作拖垮主流程。 假设扣库存 20ms、发短信
500ms、发邮件 800ms、写日志 30ms,那么一个下单接口的响应时间就是
1350ms。其中 1330ms
都耗在「用户其实不关心」的短信和邮件上。用户只想尽快看到「下单成功」。
问题二:服务强耦合。
订单服务直接调用库存服务的接口、通知服务的接口。一旦库存服务挂了或升级,订单服务跟着报错;一旦通知服务接口改了签名,订单服务必须跟着改。系统之间互相拖累,牵一发动全身。
问题三:无法削峰。 双十一零点,每秒 1
万笔下单,如果订单服务直接把这 1
万次请求同步打到数据库,数据库瞬时压力过大就会雪崩。而实际上,库存扣减、发通知这些操作,晚几秒甚至晚几十秒处理,用户完全无感。
MQ 的解法是把这三个问题一次性解决:
flowchart LR
A[订单服务<br/>接收下单请求] -->|1. 落库 + 发消息<br/>耗时约 50ms| B[(数据库)]
A -->|发消息| C[消息队列]
C -->|订阅| D[库存服务<br/>异步扣库存]
C -->|订阅| E[通知服务<br/>异步发短信邮件]
- 异步:订单服务只做「落库 + 发消息」两件快事(约
50ms),发短信、发邮件交给通知服务后台慢慢做。 - 解耦:订单服务不再直接依赖库存服务、通知服务,它只认识
MQ。库存服务改了、挂了、甚至被替换掉,订单服务一行代码都不用动。 - 削峰:1 万次下单进入 MQ
排队,消费端按自己(比如每秒 1000
条)的节奏慢慢消费,数据库压力被抹平。
一句话总结:MQ
把「同步的直接调用」变成「异步的间接投递」,换取的是速度、解耦、抗洪峰。
1.2 消息队列的两种基本模型
消息队列不管有多少种产品(RabbitMQ、Kafka、RocketMQ、ActiveMQ),其消息模型归根结底是两种的组合:
(1)点对点模型(Point-to-Point,队列 Queue)
一条消息只能被一个消费者消费,消费完就从队列里移除。典型场景:下单后「扣库存」这个动作,只能由一个库存服务实例执行一次,不能被多个实例各扣一遍。
flowchart LR
P[生产者] --> Q[队列 Queue]
Q --> C1[消费者 A]
Q --> C2[消费者 B]
(2)发布/订阅模型(Publish/Subscribe,主题
Topic)
一条消息会被广播给所有订阅了该主题的消费者。典型场景:一笔订单创建后,「库存服务要扣库存」「通知服务要发短信」「积分服务要加积分」三件事都要做,每个服务各自订阅、各收到一份消息。
flowchart LR
P[生产者] --> T[主题 Topic]
T --> C1[消费者 A]
T --> C2[消费者 B]
T --> C3[消费者 C]
实际产品中,RabbitMQ 的「交换机 +
绑定」机制可以同时表达这两种模型(fanout 交换机 = 发布订阅,direct/topic
+ 单队列 = 点对点);Kafka 的 Consumer Group
也把这两种模型统一了起来:同一个 group
内是点对点(一条消息只被组内一个实例消费),不同 group
之间是发布订阅(每个组都收到一份)。这个区别后面两章会详细展开。
1.3 三大可靠性问题(本篇主线)
这是本篇最核心的内容,也是面试必考。消息从生产者发出,到最终被消费者处理成功,中间要经过「生产者
→ Broker →
消费者」三站,任何一站出问题都会导致可靠性问题。总结起来就是三句话:
| 问题 | 含义 | 常见场景 | 解决方案(总览) |
|---|---|---|---|
| 消息丢失 | 消息该到却没到 | 生产者发出去 Broker 没收到;Broker 宕机内存数据没了;消费者收到但没处理就挂了 |
生产者 confirm 确认 + Broker 持久化 + 消费者手动 ACK |
| 消息重复 | 同一条消息被处理多次 | 网络超时后生产者重发;消费者没 ACK 导致 Broker 重投 | 消费端幂等(唯一键去重 / Redis setNX / 数据库唯一索引) |
| 消息顺序 | 需要严格有序却乱了 | 多线程并发消费;多分区并行消费 | 单队列 / 单分区 + 顺序路由,或业务上规避顺序依赖 |
为什么「丢、重、顺序」是「三座大山」?
因为它们和「可靠投递」的目标存在天然矛盾:
- 要「不丢」,就得「重试 + 持久化 +
确认」,而「重试」必然带来「重复」的可能——所以「不丢」和「不重」要靠两端配合:发送端用确认机制保证「至少一次投递」(At-Least-Once),消费端用幂等把「至少一次」兜成「恰好一次」(Exactly-Once
的效果)。 - 要「快」,就得「多分区 /
多线程并行消费」,而并行必然破坏「全局顺序」——所以「顺序」和「吞吐」也矛盾,只能二选一或在业务层规避。
这三个问题贯穿本篇后续每一章:第 2 章 RabbitMQ 的 confirm + 手动 ACK
解决「丢」,第 3 章 Kafka 的幂等消费解决「重」,第 6
章本地消息表是「丢」的终极兜底方案。记住主线:丢→确认与持久化,重→幂等,序→单分区或业务规避。
1.4 RabbitMQ vs
Kafka:如何选型
两个都是主流
MQ,但设计初衷和擅长场景差别很大。先看对比表,再看选型结论。
| 维度 | RabbitMQ | Kafka |
|---|---|---|
| 定位 | 通用的「消息代理」,AMQP 标准实现 | 分布式「流平台」,为日志、事件流而生 |
| 吞吐量 | 中等(万级/秒) | 极高(百万级/秒) |
| 延迟 | 低(毫秒级) | 低(批量写入,略高) |
| 消息回溯 | 不支持(消费即删) | 支持(按 offset 重新消费) |
| 消息顺序 | 单队列内有序 | 单分区内有序 |
| 路由能力 | 强(交换机 + 路由键 + 通配符) | 弱(只能按 Topic,分区内无路由) |
| 延迟任务 | 原生支持(TTL + 死信队列) | 不原生支持(需外部方案) |
| 运维复杂度 | 低(单节点即可跑) | 高(集群、分区、副本、ZooKeeper/KRaft) |
| 典型场景 | 业务消息、任务分发、延迟关单 | 日志采集、流处理、埋点、大数据量 |
选型口诀:
- 做业务消息(下单、扣款、发通知、延迟任务、复杂路由)→
选
RabbitMQ。它路由灵活、延迟任务原生支持、运维简单,业务语义清晰。 - 做日志 / 埋点 /
流式数据(海量数据、需要回溯、需要把同一类数据喂给多个下游)→
选 Kafka。它吞吐高、能回溯、天然支持多消费者组。 - 要事务消息(本地事务与消息严格一致)→ 阿里系场景选
RocketMQ(第 6 章会讲)。
一句话:RabbitMQ 重「路由与业务语义」,Kafka
重「吞吐与回溯」。中小型业务系统,默认选 RabbitMQ
通常不会错;当数据量达到海量日志级别再引入 Kafka。
1.5
消息投递的三种语义(At-Most-Once / At-Least-Once / Exactly-Once)
理解了三大可靠性问题,再往上抽象一层,就是消息系统里经典的「三种投递语义」。它回答的是同一个问题:一条消息,消费者最多处理几次?
| 语义 | 英文 | 含义 | 怎么实现 | 代价 |
|---|---|---|---|---|
| 至多一次 | At-Most-Once | 消息最多被处理一次,但可能丢 | 发完不管(无确认)、自动 ACK | 允许丢消息 |
| 至少一次 | At-Least-Once | 消息至少被处理一次,但可能重 | 生产者确认重试 + 消费者手动 ACK | 可能重复,需幂等 |
| 恰好一次 | Exactly-Once | 消息恰好被处理一次 | 极其困难,见下文 | 实现复杂、性能差 |
为什么说「恰好一次」在分布式消息系统里几乎不可能真正实现?
因为「生产者 → Broker →
消费者」三个环节跨网络、跨进程,任何一环的网络超时都会导致「不知道对方到底成没成功」,于是只能重试,而重试就可能重复。业界公认的结论是:
- 发送端只能保证「至少一次」(At-Least-Once):通过
confirm + 重试,保证消息一定能到 Broker,但无法保证不重复发。 - 消费端用幂等把「至少一次」收敛成「恰好一次的效果」(业界叫
Exactly-Once Semantics
的效果):即消息可能重复投递,但业务处理结果只生效一次。
所以本篇反复强调的「生产端确认 +
消费端幂等」组合,翻译成术语就是:用 At-Least-Once 投递
+ 幂等消费,逼近 Exactly-Once
的业务效果。这句话值得背下来,它是理解所有 MQ
可靠性设计的钥匙。
补充:Kafka 提供了
enable.idempotence(生产者幂等)和「事务」来进一步减少重复,但它们解决的也是「Broker
内部」的重复,跨进程的「业务处理恰好一次」最终仍要靠消费端幂等。所以无论用什么
MQ,幂等消费都是绕不开的基本功。
本章小结
MQ 用「异步投递」换来速度、解耦、削峰三大收益,但代价是引入了「丢 /
重 / 顺序」三大可靠性问题。选型上,业务消息用 RabbitMQ、海量日志流用
Kafka。记住「丢→确认与持久化,重→幂等,序→单分区」这条主线,它是贯穿全篇的钥匙。
第 2 章 RabbitMQ 实战
2.1
核心概念:交换机、队列、绑定、路由键
RabbitMQ 遵循 AMQP
协议,它的消息流转比「直接往队列里塞消息」多了一层「交换机(Exchange)」。理解这四个概念的关系,是掌握
RabbitMQ 的钥匙:
flowchart LR
P[生产者<br/>Producer] -->|消息 + 路由键 routingKey| E[交换机<br/>Exchange]
E -->|绑定关系 Binding| Q1[队列 Queue A]
E -->|绑定关系 Binding| Q2[队列 Queue B]
Q1 --> C1[消费者 Consumer]
Q2 --> C2[消费者 Consumer]
- 交换机(Exchange):消息的「邮局分拣中心」。生产者不直接把消息发给队列,而是发给交换机,并附带一个「路由键(Routing
Key)」。 - 绑定(Binding):交换机与队列之间的一条「连线」,可以理解为「分拣规则」。绑定上可以指定一个
binding key。 - 路由键(Routing
Key):生产者在发消息时指定的标签,交换机根据它和 binding key
的匹配规则,决定把消息投递到哪些队列。 - 队列(Queue):真正存储消息的容器,消费者从队列取消息。
为什么要多一层交换机?因为它把「消息投递给谁」的逻辑从生产者身上剥离了出来。生产者只负责说「我发了一条
routingKey 为 X
的消息」,至于投给哪个队列、是否广播给多个队列,完全由交换机 +
绑定的拓扑决定。新增一个消费者,只需加一条绑定,生产者代码零改动——这就是解耦的体现。
2.2 四种交换机类型
RabbitMQ 内置四类交换机,最常用的是前三种:
| 类型 | 路由规则 | 一句话理解 | 典型场景 |
|---|---|---|---|
| direct | routingKey 精确等于 binding key | 「点对点」精确匹配 | 指定某类业务消息发给指定队列 |
| fanout | 忽略 routingKey,广播给所有绑定队列 | 「大喇叭」群发 | 一条订单消息同时通知库存、通知、积分服务 |
| topic | routingKey 与 binding key 做通配符匹配( * 匹配一段,#匹配零或多段) |
「按模式」匹配 | 按 order.created、order.paid.*等模式路由 |
| headers | 按消息的 header 键值匹配,几乎不用 | 冷门 | 极少使用 |
先用一个「最小可用」配置把三种交换机跑起来,直观感受它们的区别。这个例子没有业务逻辑,纯粹演示路由行为,请重点关注
Binding 的写法。
package com.example.demo.config;
import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
* 三种交换机的最小演示配置。
* 目的:让读者直观理解 direct / fanout / topic 的路由差异。
*/
@Configuration
public class RabbitExchangeDemoConfig {
public static final String DIRECT_EXCHANGE = "demo.direct";
public static final String FANOUT_EXCHANGE = "demo.fanout";
public static final String TOPIC_EXCHANGE = "demo.topic";
// ---------- direct:精确匹配 ----------
@Bean
public DirectExchange directExchange() {
// durable=true 表示交换机持久化,Broker 重启后仍存在
return ExchangeBuilder.directExchange(DIRECT_EXCHANGE).durable(true).build();
}
@Bean
public Queue directQueueA() {
// 队列也持久化,避免重启丢队列
return QueueBuilder.durable("demo.direct.queue.a").build();
}
@Bean
public Binding directBindingA() {
// 只有 routingKey 精确等于 "order.created" 的消息才会进入该队列
return BindingBuilder.bind(directQueueA()).to(directExchange()).with("order.created");
}
// ---------- fanout:广播 ----------
@Bean
public FanoutExchange fanoutExchange() {
return ExchangeBuilder.fanoutExchange(FANOUT_EXCHANGE).durable(true).build();
}
@Bean
public Queue fanoutQueueA() {
return QueueBuilder.durable("demo.fanout.queue.a").build();
}
@Bean
public Queue fanoutQueueB() {
return QueueBuilder.durable("demo.fanout.queue.b").build();
}
@Bean
public Binding fanoutBindingA() {
// fanout 忽略 routingKey,所有绑定队列都会收到消息
return BindingBuilder.bind(fanoutQueueA()).to(fanoutExchange());
}
@Bean
public Binding fanoutBindingB() {
return BindingBuilder.bind(fanoutQueueB()).to(fanoutExchange());
}
// ---------- topic:通配符匹配 ----------
@Bean
public TopicExchange topicExchange() {
return ExchangeBuilder.topicExchange(TOPIC_EXCHANGE).durable(true).build();
}
@Bean
public Queue topicQueueOrder() {
return QueueBuilder.durable("demo.topic.queue.order").build();
}
@Bean
public Queue topicQueueAll() {
return QueueBuilder.durable("demo.topic.queue.all").build();
}
@Bean
public Binding topicBindingOrder() {
// "order.*" 匹配 order.created、order.paid 等两段式 routingKey
return BindingBuilder.bind(topicQueueOrder()).to(topicExchange()).with("order.*");
}
@Bean
public Binding topicBindingAll() {
// "order.#" 匹配以 order. 开头的任意多段 routingKey(含 order.created.sub)
return BindingBuilder.bind(topicQueueAll()).to(topicExchange()).with("order.#");
}
}
*与#的区别:*
恰好匹配一段(不含点号),#
匹配零段或多段。所以order.*能匹配
order.created,但匹配不了
order.created.sub;而order.#
两者都能匹配。
2.3
生产级配置:开启确认、持久化、手动 ACK
生产环境里,光能发消息远远不够,必须解决「消息丢失」问题。RabbitMQ
的可靠性靠三层机制叠加(对应第 1 章的「丢」):
- 生产者确认(Publisher
Confirm):生产者发消息后,Broker
回执告诉它「我收到了」,收不到回执就知道可能丢了; - Broker
持久化(durable):交换机、队列、消息都持久化到磁盘,Broker
重启不丢; - 消费者手动 ACK(Manual
Ack):消费者处理成功后手动回执basicAck,Broker
才删除消息;处理失败回执basicNack,Broker
决定是否重投。
先看 application.yml 的关键配置:
spring:
rabbitmq:
host: localhost
port: 5672
username: guest
password: guest
virtual-host: /
# 生产者确认:correlated 表示每条消息都有独立的 confirm 回调
publisher-confirm-type: correlated
# 生产者退回:消息路由不到任何队列时,回调 ReturnCallback
publisher-returns: true
template:
# mandatory=true 时,路由失败的消息才会触发 ReturnCallback
mandatory: true
listener:
simple:
# 手动 ACK:消费成功后由代码显式 ack,失败 nack,避免自动 ack 造成「处理失败也删消息」
acknowledge-mode: manual
# 预取数量:每个消费者最多同时取 10 条未 ack 消息,防止消费端被冲垮
prefetch: 10
# 消费失败时的重试(配合手动 ACK 使用,见 2.6 关键点)
retry:
enabled: true
max-attempts: 3
initial-interval: 1000
multiplier: 2.0
接着是消息序列化配置。Spring Boot 默认用
SimpleMessageConverter,会把对象序列化成 Java
序列化字节流,可读性差、跨语言不友好。生产统一用 JSON:
package com.example.demo.config;
import org.springframework.amqp.rabbit.connection.ConnectionFactory;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;
import org.springframework.amqp.support.converter.MessageConverter;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
* RabbitTemplate 统一配置:使用 JSON 序列化。
* 原因:默认的 Java 序列化字节流,跨语言无法消费,且可读性差、体积大。
*/
@Configuration
public class RabbitTemplateConfig {
@Bean
public MessageConverter jacksonMessageConverter() {
return new Jackson2JsonMessageConverter();
}
@Bean
public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
RabbitTemplate template = new RabbitTemplate(connectionFactory);
// 发送端使用 JSON 转换器,消息体以 JSON 形式存储
template.setMessageConverter(jacksonMessageConverter());
return template;
}
}
2.4
生产者:发送消息 + confirm 回调 + 路由失败退回
下面是一个生产级的订单消息生产者,它做了三件事:
- 通过
RabbitTemplate发送 JSON 消息; - 通过
ConfirmCallback感知「Broker
是否收到」,收不到则告警并补偿; - 通过
ReturnsCallback
感知「消息是否路由到队列」,路由失败说明 routingKey 或绑定配错了。
package com.example.demo.mq;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;
/**
* 订单消息生产者(生产级)。
* 通过 confirm + return 双回调,把「消息到底发没发出去」这件事显式化。
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class OrderMessageProducer {
private final RabbitTemplate rabbitTemplate;
/**
* 发送「订单已创建」事件。
*
* @param event 订单事件体
* @param bizId 业务唯一 id,用于 confirm 回调定位是哪条消息出了问题
*/
public void sendOrderCreatedEvent(OrderCreatedEvent event, String bizId) {
// CorrelationData 是本次发送的「凭据」,confirm 回调里通过它关联具体消息
CorrelationData correlationData = new CorrelationData(bizId);
// 注意:RabbitTemplate 只发送,不感知结果;结果在下面两个回调里异步返回
rabbitTemplate.convertAndSend(
RabbitOrderConfig.ORDER_EXCHANGE,
RabbitOrderConfig.ORDER_CREATED_ROUTING_KEY,
event,
correlationData
);
log.info("已发送订单创建消息,bizId={}, orderId={}", bizId, event.getOrderId());
}
}
两个回调要挂在 RabbitTemplate 上,一般放在配置类的
@PostConstruct 里注册:
package com.example.demo.config;
import jakarta.annotation.PostConstruct;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.ReturnedMessage;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.context.annotation.Configuration;
/**
* 注册生产者的 confirm 与 return 回调。
* 这是「消息不丢」的第一道防线:发出去了必须知道结果。
*/
@Slf4j
@Configuration
@RequiredArgsConstructor
public class RabbitReliabilityConfig {
private final RabbitTemplate rabbitTemplate;
@PostConstruct
public void registerCallbacks() {
// 1. 确认回调:Broker 是否成功接收(到达交换机即回调,与路由无关)
rabbitTemplate.setConfirmCallback((CorrelationData data, boolean ack, String cause) -> {
if (ack) {
log.info("消息发送确认成功,bizId={}", data == null ? null : data.getId());
} else {
// 生产上这里应:告警 + 记录待补偿 + 定时重发或人工介入
log.error("消息发送确认失败!bizId={}, cause={}",
data == null ? null : data.getId(), cause);
}
});
// 2. 退回回调:交换机收到了,但路由不到任何队列(routingKey/绑定配错)
rabbitTemplate.setReturnsCallback((ReturnedMessage returned) -> {
Message msg = returned.getMessage();
log.error("消息路由失败:exchange={}, routingKey={}, replyText={}, body={}",
returned.getExchange(), returned.getRoutingKey(),
returned.getReplyText(), new String(msg.getBody()));
});
}
}
关键点:
ConfirmCallback
是「到达交换机」的确认,ReturnsCallback
是「路由到队列失败」的退回,两者关注的是不同环节。一条消息只有「confirm
成功 且 没有 return」才算真正到了队列。
2.5 消费者:手动 ACK + 幂等
消费者这边最容易犯的错误是用默认的自动
ACK(acknowledge-mode: auto)。自动 ACK
的问题是:Broker
把消息一推给消费者,消费者方法还没执行,消息就可能被标记为已消费;一旦消费者在处理中途崩溃,这条消息就永久丢失了。
正确做法是手动 ACK:处理成功才
basicAck,处理失败 basicNack
让消息进死信队列或重投。
package com.example.demo.mq;
import com.example.demo.service.StockService;
import com.rabbitmq.client.Channel;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.support.AmqpHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;
import java.io.IOException;
/**
* 订单消息消费者(生产级)。
* 手动 ACK + 幂等,是「消息不丢、不重」的关键。
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class OrderMessageConsumer {
private final StockService stockService;
@RabbitListener(queues = RabbitOrderConfig.ORDER_CREATED_QUEUE)
public void onOrderCreated(OrderCreatedEvent event, Message message, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws IOException {
log.info("收到订单创建消息,orderId={}, deliveryTag={}", event.getOrderId(), deliveryTag);
try {
// 扣库存内部必须做幂等(唯一键 / Redis setNX),保证重复投递不会重复扣
stockService.deductStock(event);
// 处理成功:手动 ack,multiple=false 表示只确认当前这一条
channel.basicAck(deliveryTag, false);
log.info("订单创建消息处理成功,orderId={}", event.getOrderId());
} catch (Exception e) {
log.error("订单创建消息处理失败,orderId={}", event.getOrderId(), e);
// 失败:不重回队列(requeue=false),消息进入死信队列等待人工/定时补偿
// 若希望立即重试,可改为 requeue=true,但需配合消费幂等 + 重试上限,否则会死循环
channel.basicNack(deliveryTag, false, false);
}
}
}
basicNack
三个参数:(deliveryTag, multiple, requeue)。multiple
是否批量确认;requeue
是否把消息重新放回队首。生产上失败通常
requeue=false(进死信),避免「坏消息无限重试阻塞队列」。
2.6 死信队列 +
TTL:实现「订单 30 分钟未支付自动关闭」
这是本章的重头戏,也是面试高频考点。先讲清「死信」和「延迟」的关系。
什么是死信(Dead Letter)?
一条消息变成死信有三种情况:
- 消息被消费者拒绝(
basicNack/
basicReject且requeue=false); - 消息在队列里超时(超过 TTL,Time To Live
存活时间); - 队列长度满了,溢出的消息被丢弃。
什么是死信队列(DLX,Dead Letter Exchange)?
给一个队列配置一个「死信交换机」,当队列里的消息变成死信时,RabbitMQ
会自动把它转发到死信交换机,再由死信交换机路由到「死信队列」。我们监听死信队列,就能拿到这些「本该被处理但没成功」的消息,做补偿或兜底。
TTL + DLX 如何实现延迟任务?
核心思路:让消息先在一个「没消费者监听」的队列里躺 30
分钟,超时后自动变成死信,被转发到死信队列,由消费者去执行「关单」。这就是用
RabbitMQ 实现延迟任务的标准套路(RabbitMQ 本身没有原生的延迟队列,直到
3.8 才有的延迟插件是另一个方案)。
flowchart LR
P[生产者<br/>发关单消息] --> DQ[延迟队列 order.delay.queue<br/>TTL=30min<br/>无消费者]
DQ -->|30min 后超时变死信| DLX[死信交换机 order.dlx]
DLX --> DLQ[死信队列 order.dlx.queue]
DLQ --> C[关单消费者<br/>检查是否已支付<br/>未支付则关闭订单]
下面是完整配置代码:
package com.example.demo.config;
import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
* 订单业务队列 + 延迟关单队列(死信队列模式)。
* 主线:三大可靠性问题之「消息丢失」——通过持久化 + 手动 ACK + 死信兜底解决。
*/
@Configuration
public class RabbitOrderConfig {
// ---------- 正常业务:订单创建 ----------
public static final String ORDER_EXCHANGE = "order.exchange";
public static final String ORDER_CREATED_QUEUE = "order.created.queue";
public static final String ORDER_CREATED_ROUTING_KEY = "order.created";
// ---------- 延迟关单:30 分钟未支付自动关闭 ----------
public static final String ORDER_DELAY_QUEUE = "order.delay.queue"; // 延迟队列(无消费者)
public static final String ORDER_DLX_EXCHANGE = "order.dlx"; // 死信交换机
public static final String ORDER_DLX_QUEUE = "order.dlx.queue"; // 死信队列(有消费者)
public static final String ORDER_DLX_ROUTING_KEY = "order.close";
/** 订单创建交换机 */
@Bean
public DirectExchange orderExchange() {
return ExchangeBuilder.directExchange(ORDER_EXCHANGE).durable(true).build();
}
/** 订单创建队列:供库存/通知服务消费 */
@Bean
public Queue orderCreatedQueue() {
return QueueBuilder.durable(ORDER_CREATED_QUEUE).build();
}
@Bean
public Binding orderCreatedBinding() {
return BindingBuilder.bind(orderCreatedQueue())
.to(orderExchange()).with(ORDER_CREATED_ROUTING_KEY);
}
/** 延迟队列:TTL 30 分钟,且不绑定任何消费者,超时后消息自动变死信 */
@Bean
public Queue orderDelayQueue() {
return QueueBuilder.durable(ORDER_DELAY_QUEUE)
// 消息在队列里的最长存活时间:30 分钟,超时即成为死信
.ttl(30 * 60 * 1000)
// 指定死信交换机与死信路由键,死信会按此路由
.deadLetterExchange(ORDER_DLX_EXCHANGE)
.deadLetterRoutingKey(ORDER_DLX_ROUTING_KEY)
.build();
}
/** 死信交换机 */
@Bean
public DirectExchange orderDlxExchange() {
return ExchangeBuilder.directExchange(ORDER_DLX_EXCHANGE).durable(true).build();
}
/** 死信队列:真正执行「关单」的消费者监听这里 */
@Bean
public Queue orderDlxQueue() {
return QueueBuilder.durable(ORDER_DLX_QUEUE).build();
}
@Bean
public Binding orderDlxBinding() {
return BindingBuilder.bind(orderDlxQueue())
.to(orderDlxExchange()).with(ORDER_DLX_ROUTING_KEY);
}
}
关单消费者监听死信队列,收到消息后先查订单状态再决定是否关单(这就是幂等
+ 兜底的双保险):
package com.example.demo.mq;
import com.example.demo.entity.Order;
import com.example.demo.service.OrderService;
import com.rabbitmq.client.Channel;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.support.AmqpHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;
import java.io.IOException;
/**
* 延迟关单消费者:监听死信队列。
* 消息 30 分钟后才到这里,必须重新查库判断订单是否已支付——这就是「兜底」。
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class OrderCloseConsumer {
private final OrderService orderService;
@RabbitListener(queues = RabbitOrderConfig.ORDER_DLX_QUEUE)
public void onOrderClose(OrderCloseMessage message, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws IOException {
log.info("收到延迟关单消息,orderId={}", message.getOrderId());
try {
// 幂等兜底:可能用户刚好在 30 分钟内完成了支付,此时不应关单
Order order = orderService.getById(message.getOrderId());
if (order != null && "UNPAID".equals(order.getStatus())) {
orderService.closeOrder(message.getOrderId());
log.info("订单已超时关闭,orderId={}", message.getOrderId());
} else {
log.info("订单状态为 {},无需关闭,orderId={}", order == null ? "NULL" : order.getStatus(),
message.getOrderId());
}
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
log.error("关单处理异常,orderId={}", message.getOrderId(), e);
// 关单失败不重回队列(会立即再次变死信循环),记录日志人工补偿
channel.basicNack(deliveryTag, false, false);
}
}
}
在订单创建成功后,同时发两条消息:一条通知库存/通知服务,一条进延迟队列做
30 分钟后关单:
// 下单服务里:创建订单成功后发两条消息
orderMessageProducer.sendOrderCreatedEvent(event, String.valueOf(orderId)); // 即时:扣库存/发通知
orderMessageProducer.sendOrderCloseMessage(new OrderCloseMessage(orderId), "close:" + orderId); // 延迟:30 分钟后关单
2.7 关键点提示
- 手动 ACK 别漏
basicAck:忘了
ack,消息会一直「未确认」占用 prefetch
额度,最终消费者收不到新消息,队列堆积。 - 失败重试要设上限:配置里
retry.max-attempts=3,重试 3 次仍失败再走
basicNack
进死信。否则一条坏消息会无限重试,拖垮消费端。 - 延迟队列千万别给它挂消费者:一旦有人消费
order.delay.queue,消息没等到 TTL
就被拿走了,延迟失效。 - 死信兜底不等于业务兜底:关单消费者里必须重新查库判断订单状态,因为「延迟的
30 分钟里用户可能已支付」。 - 队列 TTL vs 消息 TTL:队列 TTL
是对整条队列统一生效(本文用法);消息 TTL
是逐条设置,可做不同延迟,但注意「队列里多条不同 TTL
的消息,只有队头消息超时才会触发死信」,有坑,简单场景用队列 TTL
更稳。
2.8 消费者并发度与性能调优
消息不丢、不重是「正确性」,但生产还要「快」。RabbitMQ
消费者吞吐由两个参数决定,都是高频调优点:
(1)prefetch(预取数量)
prefetch 控制「每个消费者同时最多持有多少条未 ACK
的消息」。它本质是一个限流阀:
- prefetch 太小(如
1):消费者一条一条处理,慢,但最安全(失败影响面小); - prefetch 太大(如
500):消费者一次拉一大把,快,但一旦某条失败,后续消息都要等,且单消费者宕机会有大量消息重回队列。
生产经验值:按「单条消息处理耗时 × 目标
QPS」估算。比如单条耗时 50ms,想达到 200 QPS,则单消费者
prefetch ≈ 200 × 0.05 = 10,再配合多个消费者实例水平扩容。
(2)concurrency(消费者线程数)
一个 @RabbitListener
默认只有一个消费线程。要提升并发,要么设
concurrency,要么多起几个实例:
@RabbitListener(queues = RabbitOrderConfig.ORDER_CREATED_QUEUE,
concurrency = "4") // 4 个消费者线程并发处理
public void onOrderCreated(...) { ... }
spring:
rabbitmq:
listener:
simple:
concurrency: 4 # 初始 4 个消费者线程
max-concurrency: 10 # 高峰最多扩到 10 个
调优心法:先定单条耗时 → 定单消费者 prefetch →
定目标 QPS → 算需要的消费者线程数 /
实例数。别一上来就堆线程,线程多了上下文切换开销反而拖慢;同时记住「顺序」约束——并发消费会破坏单队列顺序,要保序就别开并发(呼应第
1 章「序→单分区/单队列」)。
本章小结
RabbitMQ 通过「交换机 + 绑定 + 路由键」实现灵活路由;可靠性靠「生产者
confirm + Broker 持久化 + 消费者手动 ACK」三层叠加;「TTL +
死信队列」是它实现延迟任务的标准套路,也是面试高频考点。记住:消费端必须幂等、失败进死信、死信兜底要重查业务状态。
第 3 章 Kafka 实战
3.1
核心概念:Topic、Partition、Consumer Group、Offset
如果说 RabbitMQ 是「精密的邮局」(交换机、绑定、路由键层层分拣),那
Kafka
就是「高速的流水线日志」——它把消息当作只追加的日志,追求极致的吞吐和可回溯。理解
Kafka,先吃透四个概念:
- Topic(主题):消息的逻辑分类,类似数据库的一张表。生产者和消费者都围绕
Topic 打交道。 - Partition(分区):一个 Topic
物理上被切分成多个分区,分区是 Kafka
并行处理和存储的基本单位。每个分区内部消息严格有序(按写入顺序),但跨分区不保证全局顺序。 - Consumer
Group(消费组):一组消费者的集合。同一个组内,一个分区只能被组内一个消费者消费(避免重复);不同组之间相互独立,都能消费到完整的一份消息。这正是第
1 章说的「组内点对点、组间发布订阅」。 - Offset(偏移量):每条消息在分区内的位置编号。消费者靠记录「我消费到第几个
offset 了」来实现断点续传和回溯。
用一张图串起来:
flowchart TB
subgraph Topic[Topic: order-events]
P0[Partition 0<br/>msg0 msg1 msg2 ...]
P1[Partition 1<br/>msg0 msg1 msg2 ...]
P2[Partition 2<br/>msg0 msg1 msg2 ...]
end
P[生产者] -->|按 key 哈希路由| Topic
P0 --> CG1[Consumer Group A<br/>消费 offset 记录各自独立]
P1 --> CG1
P2 --> CG1
P0 -.-> CG2[Consumer Group B<br/>独立再消费一份]
P1 -.-> CG2
P2 -.-> CG2
和 RabbitMQ 的本质区别:RabbitMQ
的队列消费完一条就删一条(无回溯);Kafka
的消费状态由消费者自己记录的 offset 决定,Broker
不删消息(按保留时间 retention 清理)。所以 Kafka 支持「把
offset 拨回去,重新消费一遍历史数据」。
3.2 生产级配置:手动提交 +
幂等消费
Kafka 的可靠性同样落到第 1 章的「丢 / 重」上:
- 不丢:生产者
acks=all(所有副本都写入才认为成功)+
重试;消费者手动提交
offset(处理成功才提交,崩溃后从上次提交处重读,不丢)。 - 不重:正因为「手动提交 +
失败重读」会带来重复,消费端必须幂等。
先看 application.yml:
spring:
kafka:
bootstrap-servers: localhost:9092
producer:
# acks=all:消息写入所有 ISR 副本才算成功,配合重试保证不丢
acks: all
retries: 3
# 发送失败重试时,若开启幂等可避免重试产生重复(spring-kafka 默认开启 enable-idempotence)
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
consumer:
group-id: order-group
# 关闭自动提交,改由代码显式 ack,避免「还没处理完就提交」导致丢消息
enable-auto-commit: false
# 新消费组从最早消息开始读(生产可按需改 latest)
auto-offset-reset: earliest
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
properties:
# 反序列化时信任的包,避免 ClassCastException
spring.json.trusted.packages: "com.example.demo.*"
listener:
# manual_immediate:方法处理完成后立即手动提交 offset
ack-mode: manual_immediate
3.3 生产者:发送订单事件到
Kafka
Kafka 生产者比 RabbitMQ 简单,不需要配置交换机/队列(Topic
可自动创建或手动创建)。这里用 KafkaTemplate 发送,并利用
key 做分区路由。
package com.example.demo.mq;
import com.example.demo.entity.OrderCreatedEvent;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.stereotype.Component;
import java.util.concurrent.CompletableFuture;
/**
* Kafka 订单事件生产者。
* key = 订单号,保证同一订单的消息路由到同一分区(保序)。
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class KafkaOrderProducer {
public static final String TOPIC_ORDER_EVENTS = "order-events";
private final KafkaTemplate<String, OrderCreatedEvent> kafkaTemplate;
public void sendOrderEvent(OrderCreatedEvent event) {
// 用 orderId 作 key:同一订单进同一分区,分区内严格有序
String key = String.valueOf(event.getOrderId());
CompletableFuture<SendResult<String, OrderCreatedEvent>> future =
kafkaTemplate.send(TOPIC_ORDER_EVENTS, key, event);
// 异步感知发送结果:失败记日志,配合 acks=all 保证不丢
future.whenComplete((result, ex) -> {
if (ex != null) {
log.error("Kafka 发送失败,orderId={}", event.getOrderId(), ex);
} else {
log.info("Kafka 发送成功,orderId={}, partition={}, offset={}",
event.getOrderId(),
result.getRecordMetadata().partition(),
result.getRecordMetadata().offset());
}
});
}
}
3.4 消费者:手动提交 offset +
幂等
消费者是 Kafka
可靠性的关键。手动提交意味着「只有业务处理成功才提交
offset」,处理失败不提交,下次会从上次提交处重读(带来重复,所以必须幂等)。
package com.example.demo.mq;
import com.example.demo.service.OrderEventService;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Component;
/**
* Kafka 订单事件消费者(生产级)。
* 手动提交 offset + 幂等:失败不提交,重读后靠幂等保证不重复处理。
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class KafkaOrderConsumer {
private final OrderEventService orderEventService;
@KafkaListener(topics = KafkaOrderProducer.TOPIC_ORDER_EVENTS, groupId = "order-group")
public void onOrderEvent(OrderCreatedEvent event, Acknowledgment ack) {
log.info("Kafka 收到订单事件,orderId={}", event.getOrderId());
try {
// 处理内部必须幂等(唯一键 / Redis 去重):因为失败重试会重复收到同一条消息
orderEventService.process(event);
// 处理成功才提交 offset,崩溃后从该处重读,不丢
ack.acknowledge();
} catch (Exception e) {
log.error("Kafka 订单事件处理失败,orderId={},不提交 offset,等待重试", event.getOrderId(), e);
// 不调用 acknowledge():offset 不前进,稍后重读该消息
// 注意:生产上要设重试上限,坏消息进死信/告警,否则会无限重试
}
}
}
3.5 幂等消费的实现(Redis
setNX 去重)
上一节反复强调「必须幂等」,这里给出具体实现。幂等的本质是:给每条消息一个业务唯一标识,处理前先判断是否已处理过。最常用的三种方式:
| 方式 | 思路 | 适用 |
|---|---|---|
| 数据库唯一索引 | 处理结果表加 UNIQUE KEY(biz_id),插入冲突即已处理 |
有落库结果,强一致 |
| Redis setNX | setIfAbsent(bizId, 1),返回 false 即已处理 |
高并发,可设过期时间 |
| 状态机判断 | 业务状态流转(如 UNPAID → PAID 只能走一次) |
业务本身有状态 |
下面用 Redis setNX 演示(配合阶段 4 的
StringRedisTemplate):
package com.example.demo.service.impl;
import com.example.demo.entity.OrderCreatedEvent;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Service;
import java.time.Duration;
/**
* 幂等处理的通用模式:Redis setNX 去重。
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class OrderEventServiceImpl implements OrderEventService {
private final StringRedisTemplate redisTemplate;
@Override
public void process(OrderCreatedEvent event) {
// 业务唯一标识:这里用「事件类型 + 订单号」,也可用消息自带的消息 id
String idempotentKey = "idempotent:order:" + event.getOrderId();
// setIfAbsent 返回 true 表示「第一次抢到」,返回 false 表示「已处理过,跳过」
Boolean firstTime = redisTemplate.opsForValue()
.setIfAbsent(idempotentKey, "1", Duration.ofHours(24));
if (!Boolean.TRUE.equals(firstTime)) {
log.warn("重复消息,已跳过,orderId={}", event.getOrderId());
return; // 幂等:直接返回,不再重复处理
}
// 真正的业务处理(扣库存、加积分等)
doBusiness(event);
}
private void doBusiness(OrderCreatedEvent event) {
// 实际业务逻辑,这里仅示意
log.info("处理订单事件成功,orderId={}", event.getOrderId());
}
}
关键点:setNX 的 key 过期时间(这里
24h)要大于「消息可能重复投递的时间窗口」,否则过期后重复消息又会「第一次抢到」,幂等失效。若业务必须永久幂等,改用数据库唯一索引更稳妥。
3.6 关键点提示
ack-mode与enable-auto-commit
要配套:想手动提交,就必须
enable-auto-commit=false且ack-mode=manual或
manual_immediate,否则自动提交会抢先提交 offset。- 手动提交才有幂等问题的完整语境:手动提交 +
处理失败不提交 → 消息重读 →
重复,所以「手动提交」和「幂等」是一对必须同时出现的搭档。 auto-offset-reset的含义:只在「当前
group 还没有提交过 offset」时生效,earliest
从最早读,latest从最新读。已有 offset 的 group
不受它影响。- 分区保序的前提:只有「同一 key
路由到同一分区」才保序;不指定 key 时 Kafka
默认轮询分区,顺序无法保证。若业务强依赖顺序,要么用单分区,要么用 key
路由 + 单线程消费该分区。
本章小结
Kafka 靠 Partition 并行、Consumer Group 复用、Offset
回溯实现高吞吐与灵活消费。可靠性上,生产者 acks=all
保不丢,消费者「手动提交 offset +
幂等」保「不丢且不重」。记住:手动提交和幂等必须成对出现。
第 4 章 异步任务 @Async
4.1 为什么需要 @Async
有些操作「必须在请求里完成」,比如下单时校验库存是否足够、扣减库存——这些是主流程,必须同步。但另一些操作「用户不关心何时完成」,比如下单成功后发通知短信、发欢迎邮件、写操作日志——这些是旁路逻辑,理想状态是「主流程返回后,后台慢慢做」。
@Async 就是 Spring
提供的「把方法异步化」的注解:被它标注的方法会在另一个线程里执行,调用方(主线程)不等它执行完就继续往下走。
一个最小示例:
@Slf4j
@Service
public class NotifyService {
@Async
public void sendEmail(String to, String content) {
// 模拟耗时操作:发邮件
log.info("发邮件线程:{}", Thread.currentThread().getName());
}
}
但这个最小示例在生产上是有害的——它用了 Spring 默认的
SimpleAsyncTaskExecutor,每来一个任务就新建一个线程,线程用完不回收。并发一上来,线程爆炸,直接
OOM。所以 @Async 的正确打开方式是「自定义线程池 +
显式指定」,下面详讲。
4.2 生产级:自定义线程池 +
@Async
先配置线程池。线程池参数是面试高频,这里把每个参数的含义和取值逻辑讲透:
package com.example.demo.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.annotation.EnableAsync;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import java.util.concurrent.ThreadPoolExecutor;
/**
* 异步任务配置:自定义线程池,禁止用默认 SimpleAsyncTaskExecutor。
*/
@Configuration
@EnableAsync // 开启 @Async 支持(缺了这个注解 @Async 全部失效)
public class AsyncConfig {
@Bean("asyncExecutor")
public ThreadPoolTaskExecutor asyncExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
// 核心线程数:常驻线程,任务不多时也保留这些线程待命
executor.setCorePoolSize(8);
// 最大线程数:核心线程忙不过来、队列也满了时,最多扩到这个数
executor.setMaxPoolSize(16);
// 队列容量:核心线程都忙时,新任务先进入队列排队等待
executor.setQueueCapacity(500);
// 空闲线程存活时间:超过核心线程数的线程,空闲这么久后被回收
executor.setKeepAliveSeconds(60);
// 线程名前缀:方便在日志/堆栈里定位是哪类任务
executor.setThreadNamePrefix("async-");
// 拒绝策略:队列满了且线程也满了时,新任务交给调用线程同步执行(不丢任务)
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
// 等待任务都完成后才关闭线程池(配合 Spring 优雅停机)
executor.setWaitForTasksToCompleteOnShutdown(true);
executor.setAwaitTerminationSeconds(30);
// 必须调用 initialize(),否则线程池不生效
executor.initialize();
return executor;
}
}
线程池的工作流程(面试必问,和参数强相关):
- 任务来了,先看核心线程(8
个)有没有空闲:有空闲就交给核心线程执行; - 核心线程都忙,任务进队列(容量 500)排队;
- 队列也满了,再看能否扩容到最大线程数(16
个); - 最大线程也忙、队列也满,触发拒绝策略。
四种拒绝策略对比:
| 拒绝策略 | 行为 | 适用 |
|---|---|---|
AbortPolicy(默认) |
抛 RejectedExecutionException |
不能容忍丢任务,但要有人捕获异常 |
CallerRunsPolicy |
交给调用线程同步执行 | 生产推荐:丢任务风险最低,还能反压上游 |
DiscardPolicy |
静默丢弃 | 可容忍丢的日志类任务 |
DiscardOldestPolicy |
丢弃队列里最老的任务 | 优先保最新任务 |
然后在 @Async
上显式指定线程池名字(这是很多人忽略的坑,见
4.4):
package com.example.demo.service;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Service;
/**
* 通知服务:异步发短信、发邮件。
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class NotifyService {
/** 显式指定线程池,避免落到默认 SimpleAsyncTaskExecutor */
@Async("asyncExecutor")
public void sendEmail(String to, String content) {
log.info("开始发邮件,to={}, 线程={}", to, Thread.currentThread().getName());
try {
// 模拟发送邮件(生产上是调用邮件网关 SDK)
Thread.sleep(2000);
log.info("邮件发送成功,to={}", to);
} catch (InterruptedException e) {
Thread.currentThread().interrupt(); // 恢复中断标记,遵循线程中断规范
log.error("邮件发送被中断,to={}", to);
}
}
}
调用方示例:
@Service
@RequiredArgsConstructor
public class OrderServiceImpl implements OrderService {
private final NotifyService notifyService;
@Override
public Long createOrder(OrderCreateDTO dto) {
// ... 下单主流程(落库等)...
Long orderId = 123L;
// 异步发通知:主线程不等待,立即返回
notifyService.sendEmail(dto.getEmail(), "您的订单 " + orderId + " 已创建");
log.info("下单完成,orderId={},通知已异步发出", orderId);
return orderId;
}
}
4.3
需要返回值:CompletableFuture
@Async 方法返回 void
最常用,但有时异步任务需要把结果交回给主线程(比如异步查询后汇总)。这时用
CompletableFuture:
@Async("asyncExecutor")
public CompletableFuture<String> fetchUserName(Long userId) {
String name = userMapper.selectById(userId).getName();
return CompletableFuture.completedFuture(name);
}
// 调用方:可以先做别的事,最后统一取结果
CompletableFuture<String> f1 = userService.fetchUserName(1L);
CompletableFuture<String> f2 = userService.fetchUserName(2L);
String name1 = f1.get(); // 阻塞等待结果
String name2 = f2.get();
注意:
@Async返回CompletableFuture
时,方法内部不要抛受检异常,要把异常包装进
future(否则调用方get()拿不到预期结果)。
4.4 @Async
失效场景(必背,与 @Transactional 同理)
这是 @Async 最经典的坑,和阶段 1 讲的
@Transactional 失效是同一个原理——Spring
AOP 代理。先记住结论,再理解原理:
失效场景一览:
| 失效场景 | 原因 | 解决 |
|---|---|---|
| 同类自调用 | this.xxx() 没走代理对象,AOP 拦截不到 |
拆成两个 Bean,或用 AopContext 获取代理 |
方法非 public |
Spring AOP 只拦截 public 方法 | 改成 public |
未开启 @EnableAsync |
没开启异步支持 | 配置类加 @EnableAsync |
| 未指定线程池 | 用了默认 SimpleAsyncTaskExecutor,线程不回收 |
@Async("线程池名") |
| 异常被内部吞掉 | 异步方法异常不抛给主线程,容易「静默失败」 | 内部 try-catch 记日志 + 告警 |
为什么「自调用」会失效? 因为
@Async(和
@Transactional、@Cacheable
一样)是基于 AOP 代理实现的:Spring
容器里注入的是「代理对象」,你从外部调用
orderService.createOrder(),其实调用的是代理,代理在方法前后做了「切到线程池执行」的增强。但在同一个类内部用
this.xxx() 调用时,this
是原始对象本身,不是代理对象,所以增强逻辑不会执行,方法就变成了普通同步调用。
@Service
public class OrderService {
@Async("asyncExecutor")
public void asyncMethod() {
log.info("线程:{}", Thread.currentThread().getName());
}
public void caller() {
// 坑:this 是原始对象,不是代理,@Async 失效,这里会同步执行
this.asyncMethod();
}
}
验证代码——用线程名判断是否真的异步了:
@RestController
@RequiredArgsConstructor
public class AsyncTestController {
private final OrderService orderService;
@GetMapping("/async/test")
public ApiResult<String> test() {
// 1. 外部调用:走代理,@Async 生效,线程名是 async-1 之类
orderService.asyncMethod();
// 2. 自调用:不走代理,@Async 失效,线程名是 http-nio-8080-exec-1(主线程)
orderService.caller();
return ApiResult.ok("查看日志中的线程名即可验证");
}
}
观察日志:asyncMethod 的线程名是
async-1(线程池线程),caller 内部调用的
asyncMethod 线程名是
http-nio-8080-exec-1(请求线程)——证明自调用没有异步化。
修复方案:把异步方法拆到另一个 Bean
里,让调用永远「跨 Bean」:
@Service
@RequiredArgsConstructor
public class OrderService {
private final AsyncTaskService asyncTaskService; // 注入另一个 Bean
public void caller() {
asyncTaskService.asyncMethod(); // 跨 Bean 调用,走代理,@Async 生效
}
}
@Service
public class AsyncTaskService {
@Async("asyncExecutor")
public void asyncMethod() {
log.info("线程:{}", Thread.currentThread().getName());
}
}
4.5 关键点提示
- 线程池参数要按业务测算:核心线程数 ≈ CPU 核数(IO
密集可放宽到 2×CPU 核数),队列容量按「高峰期任务积压量」估,拒绝策略用
CallerRunsPolicy兜底。 - 异步方法不要依赖 Request
上下文:@Async
线程不是请求线程,RequestContextHolder、ThreadLocal
里的登录态、traceId 都取不到,需要的话要显式传递(或使用
TaskDecorator透传)。 - 异步异常不会抛回调用方:
@Async
方法的异常在独立线程抛出,调用方 catch 不到,所以异步方法内部必须自己
try-catch + 记日志。 - 验证异步一定要看线程名:这是判断「到底异步没异步」的黄金标准,别靠感觉。
本章小结
@Async 让耗时旁路操作异步化,但必须「自定义线程池 +
显式指定」,并警惕「同类自调用失效」(AOP 代理原理)。线程池四大参数 +
拒绝策略是面试必考,CallerRunsPolicy
是生产推荐。异步方法要自己处理异常、传递上下文。
第 5 章 定时任务
5.1 @Scheduled
基础:三种触发方式
定时任务解决的是「在指定时间点 /
周期性地自动执行某些逻辑」的需求,比如每天凌晨 2 点生成报表、每 30
秒同步一次数据、每隔 1 分钟扫描超时订单。
Spring 提供 @Scheduled 注解,配合
@EnableScheduling 开启。它支持三种触发方式:
| 触发方式 | 含义 | 示例 |
|---|---|---|
cron |
按 cron 表达式精确到秒 | 每天凌晨 2 点:0 0 2 * * ? |
fixedRate |
从上一次开始执行起,固定间隔再执行 | 每 60 秒:fixedRate = 60000 |
fixedDelay |
从上一次执行结束起,延迟固定时间再执行 | 上次结束后 60 秒:fixedDelay = 60000 |
最小示例:
package com.example.demo.config;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.annotation.EnableScheduling;
/** 开启定时任务支持 */
@Configuration
@EnableScheduling
public class ScheduleConfig {
}
package com.example.demo.task;
import lombok.extern.slf4j.Slf4j;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
@Slf4j
@Component
public class DemoTask {
/** cron:每天凌晨 2 点执行 */
@Scheduled(cron = "0 0 2 * * ?")
public void dailyReport() {
log.info("生成每日报表");
}
/** fixedRate:每 60 秒执行(上次开始后 60s,不管上次是否结束) */
@Scheduled(fixedRate = 60000)
public void syncData() {
log.info("同步数据");
}
/** fixedDelay:上次执行结束后延迟 5 秒再执行(不会重叠) */
@Scheduled(fixedDelay = 5000)
public void cleanTempFile() {
log.info("清理临时文件");
}
}
5.2 cron 表达式详解
cron 表达式是定时任务的「时间语言」,面试和工作都高频用到。标准格式是
7 段(最后一位是年,可省略变 6 段):
秒 分 时 日 月 周 [年]
0 0 2 * * ? = 每天凌晨 2 点
各字段取值范围和特殊字符:
| 字段 | 范围 | 特殊字符 |
|---|---|---|
| 秒 | 0-59 | , - * / |
| 分 | 0-59 | , - * / |
| 时 | 0-23 | , - * / |
| 日 | 1-31 | , - * ? / L W |
| 月 | 1-12 | , - * / |
| 周 | 1-7(1=周日) | , - * ? / L # |
常用特殊字符:
*:任意值。* * * * * ?= 每秒执行。?:不指定(日和周只能有一个用
?,避免冲突)。-:范围。0 0 9-18 * * ?= 9 点到 18
点每小时整点。/:步长。0/5 * * * * ?= 每 5
秒;0 0/30 9-18 * * ?= 9 点到 18 点每 30 分钟。,:枚举。0 0 2,14 * * ?= 每天 2 点和 14
点。
几个生产高频表达式:
| 需求 | 表达式 |
|---|---|
| 每秒 | * * * * * ? |
| 每分钟 | 0 * * * * ? |
| 每小时整点 | 0 0 * * * ? |
| 每天凌晨 2 点 | 0 0 2 * * ? |
| 每周一凌晨 3 点 | 0 0 3 ? * MON |
| 每 5 分钟 | 0 0/5 * * * ? |
| 每月 1 号凌晨 1 点 | 0 0 1 1 * ? |
注意:Spring 的 cron 是 6 段(秒开头),和 Linux
crontab 的 5 段(分开头)不一样,别搞混了。
5.3
配置定时任务线程池(TaskScheduler)
一个大坑:@Scheduled
默认用单线程执行所有任务。如果你有多个定时任务,一个任务执行时间长(比如一个报表跑
10
分钟),其他所有定时任务都会被它串行阻塞,错过触发时间。生产必须配置定时任务线程池:
package com.example.demo.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.annotation.SchedulingConfigurer;
import org.springframework.scheduling.config.ScheduledTaskRegistrar;
import java.util.concurrent.Executors;
/**
* 定时任务配置:自定义调度线程池,避免多任务互相阻塞。
*/
@Configuration
@EnableScheduling
public class ScheduleConfig implements SchedulingConfigurer {
@Override
public void configureTasks(ScheduledTaskRegistrar registrar) {
// 核心线程数按「同时可能触发的定时任务数」定,这里给 4 个足够
registrar.setScheduler(Executors.newScheduledThreadPool(4));
}
}
也可以用
@Bean提供一个
TaskScheduler,效果相同。核心思想:别用默认单线程。
5.4
分布式定时任务:多实例防重复
第二个大坑(更隐蔽):生产环境服务通常是多实例部署(比如
3 个节点)。如果每个节点都跑同一个 @Scheduled
任务,那么凌晨 2 点的报表就会被生成 3
份,扣款类任务会被执行 3 次——数据错乱。
解决方案:用分布式锁保证同一时刻只有一个实例执行。最常用的是
Redis 的 setNX(setIfAbsent),配合 Lua
脚本安全释放。
package com.example.demo.task;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.data.redis.core.script.DefaultRedisScript;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import java.time.Duration;
import java.util.Collections;
import java.util.UUID;
/**
* 分布式定时任务:用 Redis 分布式锁保证多实例下只执行一次。
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class ReportTask {
private final StringRedisTemplate redisTemplate;
private static final String LOCK_KEY = "lock:dailyReport";
// 释放锁的 Lua 脚本:先比对 value 再删除,防止误删别人的锁(value 相同才删)
private static final DefaultRedisScript<Long> UNLOCK_SCRIPT = new DefaultRedisScript<>(
"if redis.call('get', KEYS[1]) == ARGV[1] then " +
" return redis.call('del', KEYS[1]) " +
"else return 0 end", Long.class);
@Scheduled(cron = "0 0 2 * * ?")
public void dailyReport() {
// 锁的唯一 value:用 UUID 区分「这把锁是谁加的」
String lockValue = UUID.randomUUID().toString();
// 抢锁:setIfAbsent 返回 true 表示抢到,false 表示其他实例已在执行
Boolean locked = redisTemplate.opsForValue()
.setIfAbsent(LOCK_KEY, lockValue, Duration.ofMinutes(30));
if (!Boolean.TRUE.equals(locked)) {
log.info("其他实例正在执行报表任务,本实例跳过");
return; // 没抢到锁,直接跳过,避免重复执行
}
try {
// 只有抢到锁的实例才真正执行任务
log.info("开始生成每日报表...");
doGenerateReport();
} finally {
// 安全释放:只有锁还是自己的才删(防止任务超时后锁过期、误删别人新加的锁)
redisTemplate.execute(UNLOCK_SCRIPT, Collections.singletonList(LOCK_KEY), lockValue);
}
}
private void doGenerateReport() {
// 实际报表生成逻辑
log.info("报表生成完成");
}
}
两个细节必须注意:
- 锁要设过期时间(这里 30
分钟),否则执行实例崩溃,锁永远不释放,任务永远被跳过;- 释放锁要用 Lua 脚本先比对 value
再删,否则任务执行超时(锁已过期被别人抢到)后,本实例的
finally会把别人的锁误删。
5.5
分布式调度框架:XXL-Job / ElasticJob
当定时任务多到需要「可视化、失败重试、分片执行、动态调度」时,手写
@Scheduled +
分布式锁就力不从心了。生产上常用分布式任务调度框架:
| 框架 | 特点 | 适用 |
|---|---|---|
| XXL-Job | 阿里开源,Web 管理界面、失败重试、动态修改 cron、分片广播 | 中小团队首选,Java 生态成熟 |
| ElasticJob | 当当开源,基于 ZooKeeper,强调分片和数据倾斜 | 大数据量分片场景 |
| Spring Cloud Task / Quartz | 前者偏批处理,后者单机强 | 特定场景 |
它们共同解决了 @Scheduled
的两个致命短板:多实例重复执行(框架内置调度中心统一触发)和任务不可观测(Web
界面看到每次执行的成功/失败、日志、重试)。本文实战部分给出 XXL-Job
的接入方向即可(完整接入见其官方文档)。
5.6 关键点提示
fixedRatevs
fixedDelay:fixedRate
是「固定频率」,任务执行时间超过间隔时会并发重叠;fixedDelay
是「固定间隔」,一定等上次结束才计时,不会重叠。要防重叠用
fixedDelay。- 默认单线程是大坑:多个
@Scheduled
必须配TaskScheduler线程池。 - 多实例必须加分布式锁:否则每个实例都执行一遍,扣款类任务会造成生产事故。
- 锁的过期时间要大于任务最坏执行时间,否则锁提前过期,任务可能被并发执行。
本章小结
@Scheduled 提供 cron / fixedRate / fixedDelay
三种触发方式,但默认单线程会导致任务互相阻塞,多实例部署会导致任务重复执行。生产上必须配调度线程池
+ 分布式锁(Redis setNX + Lua 释放),任务多了再上 XXL-Job
这类框架。
第 6 章
分布式事务(最终一致性)
6.1 为什么单库事务不够了
在阶段 4 里,你学会了用 @Transactional
保证「一个数据库里多个操作要么全成功、要么全失败」。这套机制在单体单库下完美,但一旦业务拆成多个服务、多个数据库,它就失效了。
看一个典型场景——下单:
订单服务(数据库 A):
begin
插入订单表 ✓
发送「订单已创建」消息给 MQ
commit
问题出在哪?「插入订单」和「发送消息」这两个操作无法放进同一个数据库事务里——数据库事务只能管数据库里的操作,管不了发消息。于是出现两种灾难:
- 先提交数据库,后发消息:如果 commit
之后、发消息之前,进程崩溃了 → 订单写进去了,但消息没发出去 →
库存服务永远不知道这笔订单,业务数据与消息不一致。 - 先发消息,后提交数据库:如果发消息之后、commit
之前,数据库回滚了 → 消息发出去了,但订单没写进去 →
库存服务收到一个「幽灵订单」,扣了不该扣的库存。
这就是分布式事务要解决的核心问题:跨资源(数据库 +
MQ,或数据库 A + 数据库
B)的一致性。由于分布式环境下无法做到真正的强一致(CAP
理论),业界普遍接受最终一致性:允许短暂不一致,但最终会一致。
6.2 三大方案对比
| 方案 | 核心思路 | 一致性 | 适用场景 |
|---|---|---|---|
| 本地消息表 | 业务数据 + 消息表同库同事务写入,再异步投递 MQ | 最终一致 | 可靠性要求高、改动小,最经典 |
| 事务消息(RocketMQ) | 半消息 + 事务状态回查,MQ 帮你确认本地事务结果 | 最终一致 | 阿里系、扣款等强可靠场景 |
| Seata(AT/TCC/Saga) | 全局事务协调器 + 各分支事务 | AT 强一致,TCC/Saga 最终一致 | 跨服务强一致 / 长事务 |
本篇重点讲透本地消息表(面试必考、最通用、不依赖特定
MQ),事务消息和 Seata 讲清原理和适用边界。
6.3 本地消息表(重点,完整实现)
核心思想:在同一个数据库里多建一张
local_message
表,业务代码在同一个本地事务里「写业务数据 +
写消息记录」,然后由一个定时任务扫描这张表,把「待投递」的消息异步发给
MQ,成功后更新状态。因为「业务数据」和「消息记录」在同一个数据库事务里,要么都成功、要么都失败,从根上解决了「消息丢了」的问题。
sequenceDiagram
participant S as 业务服务
participant DB as 数据库(业务表+消息表)
participant J as 投递定时任务
participant M as MQ
participant C as 消费方
S->>DB: 本地事务:写业务数据 + 写消息表(状态=待投递)
Note over DB: 同库同事务,要么都成功要么都失败
J->>DB: 定时扫描 status=待投递 的消息
J->>M: 投递消息
M-->>J: 投递成功
J->>DB: 更新消息 status=已投递
M->>C: 消费方处理 + 幂等 ACK
第一步:建消息表
CREATE TABLE `local_message` (
`id` BIGINT NOT NULL AUTO_INCREMENT COMMENT '主键',
`biz_type` VARCHAR(64) NOT NULL COMMENT '业务类型,如 ORDER_CREATED',
`biz_key` VARCHAR(128) NOT NULL COMMENT '业务唯一键,用于幂等和去重',
`payload` TEXT NOT NULL COMMENT '消息体(JSON)',
`status` TINYINT NOT NULL DEFAULT 0 COMMENT '0=待投递 1=已投递',
`retry_count` INT NOT NULL DEFAULT 0 COMMENT '已重试次数',
`create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`update_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_biz_key` (`biz_key`) -- 唯一索引:保证同一业务不重复投递
) ENGINE = InnoDB DEFAULT CHARSET = utf8mb4 COMMENT = '本地消息表';
第二步:实体 + Mapper
package com.example.demo.entity;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import java.time.LocalDateTime;
@Data
@TableName("local_message")
public class LocalMessage {
@TableId(type = IdType.AUTO)
private Long id;
/** 业务类型 */
private String bizType;
/** 业务唯一键(如订单号),用于幂等 */
private String bizKey;
/** 消息体 JSON */
private String payload;
/** 0=待投递 1=已投递 */
private Integer status;
/** 重试次数 */
private Integer retryCount;
private LocalDateTime createTime;
private LocalDateTime updateTime;
}
package com.example.demo.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.example.demo.entity.LocalMessage;
import org.apache.ibatis.annotations.Mapper;
@Mapper
public interface LocalMessageMapper extends BaseMapper<LocalMessage> {
}
第三步:业务 + 消息表同事务(核心)
这是整个模式最关键的一步:@Transactional
包住「写订单 +
写消息表」。因为两者在同一个本地事务里,数据库保证它们要么都成功、要么都回滚。
package com.example.demo.service.impl;
import com.alibaba.fastjson2.JSON;
import com.example.demo.dto.OrderCreateDTO;
import com.example.demo.entity.LocalMessage;
import com.example.demo.entity.Order;
import com.example.demo.entity.OrderCreatedEvent;
import com.example.demo.mapper.LocalMessageMapper;
import com.example.demo.mapper.OrderMapper;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
/**
* 下单服务:本地消息表模式。
* 关键:业务数据 + 消息记录在【同一个本地事务】里写入。
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class OrderServiceImpl implements OrderService {
private final OrderMapper orderMapper;
private final LocalMessageMapper messageMapper;
@Override
@Transactional(rollbackFor = Exception.class) // 关键:同库同事务
public Long createOrder(OrderCreateDTO dto) {
// 1. 写业务数据
Order order = new Order();
order.setUserId(dto.getUserId());
order.setAmount(dto.getAmount());
order.setStatus("UNPAID");
orderMapper.insert(order);
Long orderId = order.getId();
// 2. 写消息记录(与上面在同一个事务里,保证不丢)
OrderCreatedEvent event = new OrderCreatedEvent(orderId, dto.getUserId(), dto.getAmount());
LocalMessage message = new LocalMessage();
message.setBizType("ORDER_CREATED");
message.setBizKey(String.valueOf(orderId)); // 唯一键 = 订单号
message.setPayload(JSON.toJSONString(event)); // 消息体 JSON
message.setStatus(0); // 待投递
message.setRetryCount(0);
messageMapper.insert(message);
// 3. 事务提交后,两条数据同时落库;之后由定时任务投递 MQ
log.info("订单创建成功,orderId={},消息记录已同事务写入", orderId);
return orderId;
}
}
为什么这能保证「消息不丢」?
因为消息记录和业务数据同生共死:订单写成功了,消息记录一定也在(事务提交);订单写失败,消息记录也不会单独存在(事务回滚)。消息「凭空丢失」的窗口被数据库事务彻底封死了。
第四步:定时任务扫描投递
业务写完就结束了,真正的「发消息」交给定时任务去做:每隔几秒扫描
status=待投递 的记录,投递给 MQ,成功则更新为已投递。
package com.example.demo.task;
import com.example.demo.entity.LocalMessage;
import com.example.demo.mapper.LocalMessageMapper;
import com.example.demo.mq.OrderMessageProducer;
import com.alibaba.fastjson2.JSON;
import com.example.demo.entity.OrderCreatedEvent;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import java.util.List;
/**
* 本地消息表投递任务:定时扫描待投递消息,发送到 MQ。
* 这里同样要防多实例重复投递——靠消息表的唯一键 + 状态机兜底。
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class MessageRelayTask {
private final LocalMessageMapper messageMapper;
private final OrderMessageProducer orderMessageProducer;
/** 每 5 秒扫描一次待投递消息 */
@Scheduled(fixedDelay = 5000)
public void relay() {
// 只查待投递的消息,一次最多 100 条,避免一次扫太多
List<LocalMessage> messages = messageMapper.selectList(
new com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper<LocalMessage>()
.eq(LocalMessage::getStatus, 0)
.last("limit 100"));
for (LocalMessage msg : messages) {
try {
OrderCreatedEvent event = JSON.parseObject(msg.getPayload(), OrderCreatedEvent.class);
// 投递 MQ(RabbitMQ 的 confirm 机制已保证 Broker 收到)
orderMessageProducer.sendOrderCreatedEvent(event, msg.getBizKey());
// 投递成功,更新状态为已投递
msg.setStatus(1);
messageMapper.updateById(msg);
} catch (Exception e) {
log.error("消息投递失败,bizKey={}", msg.getBizKey(), e);
// 失败不更新状态,下次扫描继续重试;可结合 retry_count 做告警
}
}
}
}
投递重复怎么办?
如果投递成功后、更新状态前崩溃了,下次扫描会再投递一次,导致消费端收到重复消息。但没关系——这正是第
1 章主线「重→幂等」的用武之地:消费端靠
bizKey(订单号)做幂等,重复消息被去重。本地消息表
+ 消费端幂等,构成「消息不丢且不重」的完整闭环。
第五步:消费端幂等(复用第 3 章 / 第 2
章的幂等逻辑)
消费端拿到消息后,用 bizKey 做幂等去重(Redis setNX
或唯一索引),保证即使重复投递也只处理一次。这与第 2 章
OrderMessageConsumer、第 3 章
KafkaOrderConsumer 完全一致,不再赘述。
6.4 事务消息(RocketMQ)与
Seata(了解原理)
事务消息(RocketMQ) 是本地消息表的「MQ
内建版」,原理是「半消息 + 状态回查」:
- 生产者先发一条半消息(half
message),此时消费者看不到它; - 生产者执行本地事务(写订单);
- 本地事务成功 → 提交半消息,消费者可见;失败 →
回滚半消息,消费者永远看不到; - 如果生产者执行本地事务后没来得及提交/回滚就崩溃了,MQ
会定时回查生产者「你那个本地事务到底成没成」,据此决定提交还是回滚半消息。
它省去了本地消息表和扫描任务,但只限 RocketMQ
支持,且需要处理回查逻辑。阿里系业务首选。
Seata
是一个分布式事务中间件,通过「全局事务协调器(TC)+
各服务分支事务(TM/RM)」实现跨服务一致性,有三种模式:
| 模式 | 原理 | 特点 |
|---|---|---|
| AT | 自动生成反向 SQL(undo log),自动回滚 | 无侵入,强一致,但性能损耗大 |
| TCC | 每个服务实现 Try/Confirm/Cancel 三个方法 | 手动编码,性能好,侵入大 |
| Saga | 长事务拆成多个本地事务,失败反向补偿 | 适合长流程,无锁 |
选型建议:本地消息表能解决 90% 的「业务 +
消息」一致性问题,先掌握它;Seata
只在「跨多个服务、多个数据库、要强一致」的复杂场景才引入,学习成本和运维成本都不低。
6.5 关键点提示
- 本地消息表必须「同库同事务」:如果消息表和业务表不在同一个数据库,本地事务就罩不住它们了,模式失效。这是理解这个模式的前提。
- 投递任务要有「状态机 +
重试上限」:status是状态机(待投递 →
已投递),retry_count
用于告警和兜底,防止一条坏消息无限重试。 - 消费端幂等是最终闭环:本地消息表保证「不丢」,幂等保证「不重」,两者缺一不可。
- 消息体要存 JSON:
payload存序列化后的
JSON 字符串,投递时反序列化,方便跨语言和排查。
本章小结
单库事务管不了「数据库 +
MQ」的跨资源一致性,业界用最终一致性解决。本地消息表是最经典、最通用的方案:业务数据和消息记录同库同事务写入,定时任务扫描投递,消费端幂等兜底。事务消息(RocketMQ)是它的内建版,Seata
用于跨服务强一致场景。
5.
生产级实战项目:下单异步解耦(完整可运行)
5.1 项目目标与架构
这个实战项目把本篇 80% 以上的知识点串成一条完整链路:
- 下单接口:本地事务里「写订单 +
写本地消息表」,保证消息不丢(第 6 章); - 消息投递任务:定时扫描消息表,按业务类型投递到
RabbitMQ 的不同队列(第 6 章 + 第 2 章); - 库存消费端:手动 ACK + Redis 幂等,扣减库存(第 2
章); - 通知消费端:
@Async异步发通知邮件(第
4 章); - 延迟关单:TTL 30 分钟 +
死信队列,未支付自动关单(第 2 章); - 定时任务:Redis 分布式锁防多实例重复(第 5
章)。
POST /order/create
│
▼
OrderServiceImpl.createOrder ──@Transactional──► 数据库(orders 表 + local_message 表,同事务)
│
▼
MessageRelayTask(定时扫描 local_message)
├─ ORDER_CREATED ──────► order.exchange ──► order.created.queue ──► StockConsumer(手动ACK+幂等)
│ └─► NotifyService(@Async 发邮件)
└─ ORDER_CLOSE_DELAY ──► order.delay.queue(TTL 30min)──超时──► DLX ──► order.dlx.queue ──► CloseConsumer(重查订单状态)
5.2 目录结构
com.example.demo
├── DemoApplication.java # 启动类
├── common # 公共基础设施(与阶段 1 逐字一致)
│ ├── ApiResult.java
│ ├── ErrorCode.java
│ └── exception
│ ├── BizException.java
│ └── GlobalExceptionHandler.java
├── controller
│ └── OrderController.java # 下单接口
├── service
│ ├── OrderService.java # 接口
│ ├── StockService.java # 扣库存(幂等)
│ ├── OrderEventService.java # Kafka 幂等处理(第 3 章)
│ ├── NotifyService.java # 异步通知(@Async)
│ └── impl/
│ ├── OrderServiceImpl.java # 下单 + 本地消息表
│ ├── StockServiceImpl.java # 扣库存实现(Redis 幂等)
│ └── OrderEventServiceImpl.java # 第 3 章已给出
├── mapper
│ ├── OrderMapper.java
│ └── LocalMessageMapper.java
├── entity
│ ├── Order.java
│ ├── LocalMessage.java # 第 6 章已给出
│ ├── OrderCreatedEvent.java
│ └── OrderCloseMessage.java
├── dto
│ └── OrderCreateDTO.java
├── mq
│ ├── OrderMessageProducer.java # 生产者(confirm + 双消息)
│ ├── OrderMessageConsumer.java # 库存消费(手动 ACK + 幂等)
│ └── OrderCloseConsumer.java # 延迟关单(死信消费,第 2 章已给出)
├── task
│ ├── MessageRelayTask.java # 消息表投递任务
│ └── ReportTask.java # 分布式定时任务(第 5 章已给出)
└── config
├── RabbitOrderConfig.java # 交换机/队列/死信(第 2 章已给出)
├── RabbitTemplateConfig.java # JSON 序列化(第 2 章已给出)
├── RabbitReliabilityConfig.java # confirm/return 回调(第 2 章已给出)
├── AsyncConfig.java # 异步线程池(第 4 章已给出)
└── ScheduleConfig.java # 定时任务线程池(第 5 章已给出)
已在第 2/4/5/6
章给出完整源码的配置类、实体、任务,这里不再重复粘贴,请按目录结构对照回填。下面给出本实战项目新增或需要补齐的业务代码。
5.3 建库建表 SQL
-- 订单表(避免与 MySQL 关键字 order 冲突,用 orders)
CREATE TABLE `orders` (
`id` BIGINT NOT NULL AUTO_INCREMENT COMMENT '订单号',
`user_id` BIGINT NOT NULL COMMENT '用户 id',
`amount` DECIMAL(10,2) NOT NULL COMMENT '金额',
`status` VARCHAR(16) NOT NULL DEFAULT 'UNPAID' COMMENT '状态:UNPAID/PAID/CLOSED',
`create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
`update_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
PRIMARY KEY (`id`)
) ENGINE = InnoDB DEFAULT CHARSET = utf8mb4 COMMENT = '订单表';
-- 本地消息表(第 6 章已给出建表语句,此处省略,直接复制第 6 章 6.3 的 SQL 即可)
5.4 实体与 DTO
package com.example.demo.entity;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import java.math.BigDecimal;
import java.time.LocalDateTime;
@Data
@TableName("orders")
public class Order {
@TableId(type = IdType.AUTO)
private Long id;
private Long userId;
private BigDecimal amount;
/** UNPAID / PAID / CLOSED */
private String status;
private LocalDateTime createTime;
private LocalDateTime updateTime;
}
package com.example.demo.dto;
import jakarta.validation.constraints.DecimalMin;
import jakarta.validation.constraints.NotNull;
import lombok.Data;
import java.math.BigDecimal;
@Data
public class OrderCreateDTO {
@NotNull(message = "用户 id 不能为空")
private Long userId;
@NotNull(message = "金额不能为空")
@DecimalMin(value = "0.01", message = "金额必须大于 0")
private BigDecimal amount;
}
package com.example.demo.entity;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.math.BigDecimal;
/** 订单创建事件:消息体(JSON 序列化后存入消息表 / 投递 MQ) */
@Data
@NoArgsConstructor
@AllArgsConstructor
public class OrderCreatedEvent {
private Long orderId;
private Long userId;
private BigDecimal amount;
}
package com.example.demo.entity;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
/** 延迟关单消息体 */
@Data
@NoArgsConstructor
@AllArgsConstructor
public class OrderCloseMessage {
private Long orderId;
}
package com.example.demo.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.example.demo.entity.Order;
import org.apache.ibatis.annotations.Mapper;
@Mapper
public interface OrderMapper extends BaseMapper<Order> {
}
5.5 生产者:confirm 确认 +
发两条消息
package com.example.demo.mq;
import com.example.demo.config.RabbitOrderConfig;
import com.example.demo.entity.OrderCloseMessage;
import com.example.demo.entity.OrderCreatedEvent;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;
/**
* 订单消息生产者:统一入口,负责发「订单创建」与「延迟关单」两类消息。
* confirm / return 回调已在 RabbitReliabilityConfig 中注册。
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class OrderMessageProducer {
private final RabbitTemplate rabbitTemplate;
/** 发「订单已创建」事件(即时投递给库存/通知消费端) */
public void sendOrderCreatedEvent(OrderCreatedEvent event, String bizId) {
CorrelationData correlationData = new CorrelationData(bizId);
rabbitTemplate.convertAndSend(
RabbitOrderConfig.ORDER_EXCHANGE,
RabbitOrderConfig.ORDER_CREATED_ROUTING_KEY,
event,
correlationData);
log.info("已发送订单创建消息,orderId={}", event.getOrderId());
}
/** 发「延迟关单」消息:投递到 order.delay.queue(TTL 30 分钟,超时进死信) */
public void sendOrderCloseMessage(OrderCloseMessage message, String bizId) {
CorrelationData correlationData = new CorrelationData(bizId);
// 默认交换机(空串)+ routingKey=队列名,消息直接进入指定队列
rabbitTemplate.convertAndSend("", RabbitOrderConfig.ORDER_DELAY_QUEUE, message, correlationData);
log.info("已发送延迟关单消息,orderId={}", message.getOrderId());
}
}
5.6
下单服务:本地消息表(同库同事务)
package com.example.demo.service;
import com.example.demo.dto.OrderCreateDTO;
public interface OrderService {
Long createOrder(OrderCreateDTO dto);
Order getById(Long orderId);
void closeOrder(Long orderId);
}
package com.example.demo.service.impl;
import com.alibaba.fastjson2.JSON;
import com.example.demo.dto.OrderCreateDTO;
import com.example.demo.entity.LocalMessage;
import com.example.demo.entity.Order;
import com.example.demo.entity.OrderCloseMessage;
import com.example.demo.entity.OrderCreatedEvent;
import com.example.demo.mapper.LocalMessageMapper;
import com.example.demo.mapper.OrderMapper;
import com.example.demo.service.OrderService;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
/**
* 下单服务:本地消息表模式。
* 关键:订单 + 消息记录在【同一个本地事务】里写入,消息永不丢失。
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class OrderServiceImpl implements OrderService {
private final OrderMapper orderMapper;
private final LocalMessageMapper messageMapper;
@Override
@Transactional(rollbackFor = Exception.class)
public Long createOrder(OrderCreateDTO dto) {
// 1. 写业务数据(订单)
Order order = new Order();
order.setUserId(dto.getUserId());
order.setAmount(dto.getAmount());
order.setStatus("UNPAID");
orderMapper.insert(order);
Long orderId = order.getId();
// 2. 写「订单创建」消息记录(同事务)
OrderCreatedEvent createdEvent = new OrderCreatedEvent(orderId, dto.getUserId(), dto.getAmount());
LocalMessage createdMsg = new LocalMessage();
createdMsg.setBizType("ORDER_CREATED");
createdMsg.setBizKey("order:" + orderId);
createdMsg.setPayload(JSON.toJSONString(createdEvent));
createdMsg.setStatus(0);
createdMsg.setRetryCount(0);
messageMapper.insert(createdMsg);
// 3. 写「延迟关单」消息记录(同事务)
OrderCloseMessage closeMessage = new OrderCloseMessage(orderId);
LocalMessage closeMsg = new LocalMessage();
closeMsg.setBizType("ORDER_CLOSE_DELAY");
closeMsg.setBizKey("close:" + orderId);
closeMsg.setPayload(JSON.toJSONString(closeMessage));
closeMsg.setStatus(0);
closeMsg.setRetryCount(0);
messageMapper.insert(closeMsg);
// 4. 事务提交后,两条消息记录随订单一起落库,由 MessageRelayTask 投递
log.info("下单成功,orderId={},订单与消息记录已同事务写入", orderId);
return orderId;
}
@Override
public Order getById(Long orderId) {
return orderMapper.selectById(orderId);
}
@Override
public void closeOrder(Long orderId) {
Order order = orderMapper.selectById(orderId);
if (order != null && "UNPAID".equals(order.getStatus())) {
order.setStatus("CLOSED");
orderMapper.updateById(order);
log.info("订单已关闭,orderId={}", orderId);
}
}
}
5.7
消息表投递任务(按业务类型路由)
package com.example.demo.task;
import com.alibaba.fastjson2.JSON;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.example.demo.entity.LocalMessage;
import com.example.demo.entity.OrderCloseMessage;
import com.example.demo.entity.OrderCreatedEvent;
import com.example.demo.mapper.LocalMessageMapper;
import com.example.demo.mq.OrderMessageProducer;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import java.util.List;
/**
* 本地消息表投递任务:扫描待投递消息,按 bizType 路由到不同队列。
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class MessageRelayTask {
private final LocalMessageMapper messageMapper;
private final OrderMessageProducer orderMessageProducer;
@Scheduled(fixedDelay = 5000)
public void relay() {
List<LocalMessage> messages = messageMapper.selectList(
new LambdaQueryWrapper<LocalMessage>()
.eq(LocalMessage::getStatus, 0)
.last("limit 100"));
for (LocalMessage msg : messages) {
try {
if ("ORDER_CREATED".equals(msg.getBizType())) {
OrderCreatedEvent event = JSON.parseObject(msg.getPayload(), OrderCreatedEvent.class);
orderMessageProducer.sendOrderCreatedEvent(event, msg.getBizKey());
} else if ("ORDER_CLOSE_DELAY".equals(msg.getBizType())) {
OrderCloseMessage close = JSON.parseObject(msg.getPayload(), OrderCloseMessage.class);
orderMessageProducer.sendOrderCloseMessage(close, msg.getBizKey());
}
msg.setStatus(1); // 投递成功 → 已投递
messageMapper.updateById(msg);
} catch (Exception e) {
log.error("消息投递失败,bizKey={}, bizType={}", msg.getBizKey(), msg.getBizType(), e);
// 失败保留 status=0,下次扫描继续重试
}
}
}
}
5.8 库存消费端:手动 ACK
+ 幂等 + 异步通知
package com.example.demo.service;
import com.example.demo.entity.OrderCreatedEvent;
public interface StockService {
void deductStock(OrderCreatedEvent event);
}
package com.example.demo.service.impl;
import com.example.demo.entity.OrderCreatedEvent;
import com.example.demo.service.NotifyService;
import com.example.demo.service.StockService;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Service;
import java.time.Duration;
/**
* 扣库存服务:Redis setNX 幂等 + 异步通知。
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class StockServiceImpl implements StockService {
private final StringRedisTemplate redisTemplate;
private final NotifyService notifyService;
@Override
public void deductStock(OrderCreatedEvent event) {
// 幂等:同一订单只扣一次库存
String idempotentKey = "idempotent:stock:" + event.getOrderId();
Boolean firstTime = redisTemplate.opsForValue()
.setIfAbsent(idempotentKey, "1", Duration.ofHours(24));
if (!Boolean.TRUE.equals(firstTime)) {
log.warn("重复扣库存消息,已跳过,orderId={}", event.getOrderId());
return;
}
// 真正扣库存(生产上是 UPDATE stock SET num = num - ? WHERE ...)
log.info("扣减库存成功,orderId={}, amount={}", event.getOrderId(), event.getAmount());
// 异步发通知(@Async 不阻塞消费线程)
notifyService.sendEmail(event.getUserId(), "您的订单 " + event.getOrderId() + " 已创建");
}
}
package com.example.demo.service;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Service;
/**
* 通知服务:异步发邮件(自定义线程池,见 AsyncConfig)。
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class NotifyService {
@Async("asyncExecutor")
public void sendEmail(Long userId, String content) {
log.info("异步发通知,userId={}, content={}, 线程={}",
userId, content, Thread.currentThread().getName());
// 生产:调用邮件网关 SDK
}
}
package com.example.demo.mq;
import com.example.demo.config.RabbitOrderConfig;
import com.example.demo.entity.OrderCreatedEvent;
import com.example.demo.service.StockService;
import com.rabbitmq.client.Channel;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.support.AmqpHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;
import java.io.IOException;
/**
* 订单创建消费端:手动 ACK,处理失败进死信。
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class OrderMessageConsumer {
private final StockService stockService;
@RabbitListener(queues = RabbitOrderConfig.ORDER_CREATED_QUEUE)
public void onOrderCreated(OrderCreatedEvent event, Message message, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws IOException {
log.info("收到订单创建消息,orderId={}", event.getOrderId());
try {
stockService.deductStock(event); // 内部已做幂等
channel.basicAck(deliveryTag, false); // 成功手动 ACK
} catch (Exception e) {
log.error("订单创建消息处理失败,orderId={}", event.getOrderId(), e);
channel.basicNack(deliveryTag, false, false); // 失败进死信,不重回队列
}
}
}
5.9 下单接口
package com.example.demo.controller;
import com.example.demo.common.ApiResult;
import com.example.demo.dto.OrderCreateDTO;
import com.example.demo.service.OrderService;
import jakarta.validation.Valid;
import lombok.RequiredArgsConstructor;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
@RestController
@RequestMapping("/order")
@RequiredArgsConstructor
public class OrderController {
private final OrderService orderService;
@PostMapping("/create")
public ApiResult<Long> create(@Valid @RequestBody OrderCreateDTO dto) {
Long orderId = orderService.createOrder(dto);
return ApiResult.ok(orderId);
}
}
5.10 application.yml(完整版)
spring:
datasource:
url: jdbc:mysql://localhost:3306/demo?useUnicode=true&characterEncoding=utf8&serverTimezone=Asia/Shanghai
username: root
password: root
driver-class-name: com.mysql.cj.jdbc.Driver
data:
redis:
host: localhost
port: 6379
rabbitmq:
host: localhost
port: 5672
username: guest
password: guest
publisher-confirm-type: correlated
publisher-returns: true
template:
mandatory: true
listener:
simple:
acknowledge-mode: manual
prefetch: 10
retry:
enabled: true
max-attempts: 3
kafka:
bootstrap-servers: localhost:9092
producer:
acks: all
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
consumer:
group-id: order-group
enable-auto-commit: false
auto-offset-reset: earliest
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
properties:
spring.json.trusted.packages: "com.example.demo.*"
listener:
ack-mode: manual_immediate
mybatis-plus:
configuration:
map-underscore-to-camel-case: true
log-impl: org.apache.ibatis.logging.stdout.StdOutImpl
5.11 公共类清单(与阶段 1
逐字一致)
import lombok.Data;
import org.slf4j.MDC;
@Data
public class ApiResult<T> {
private int code; // 0=成功,非 0=错误码
private String message; // 提示信息
private T data; // 业务数据
private String traceId; // 链路追踪 id
public static <T> ApiResult<T> ok(T data) {
ApiResult<T> r = new ApiResult<>();
r.setCode(ErrorCode.SUCCESS.getCode());
r.setMessage(ErrorCode.SUCCESS.getMessage());
r.setData(data);
r.setTraceId(MDC.get("traceId"));
return r;
}
public static <T> ApiResult<T> ok() {
return ok(null);
}
public static <T> ApiResult<T> fail(int code, String message) {
ApiResult<T> r = new ApiResult<>();
r.setCode(code);
r.setMessage(message);
r.setTraceId(MDC.get("traceId"));
return r;
}
public static <T> ApiResult<T> fail(ErrorCode ec) {
return fail(ec.getCode(), ec.getMessage());
}
}
package com.example.demo.common;
public enum ErrorCode {
SUCCESS(0, "success"),
PARAM_ERROR(40001, "参数错误"),
UNAUTHORIZED(40101, "未登录或登录已过期"),
FORBIDDEN(40301, "无权限访问"),
USER_NOT_FOUND(40401, "用户不存在"),
SYSTEM_ERROR(50000, "系统繁忙,请稍后重试"),
;
private final int code;
private final String message;
ErrorCode(int code, String message) {
this.code = code;
this.message = message;
}
public int getCode() { return code; }
public String getMessage() { return message; }
}
package com.example.demo.common.exception;
import com.example.demo.common.ErrorCode;
public class BizException extends RuntimeException {
private final int code;
public BizException(ErrorCode ec) {
super(ec.getMessage());
this.code = ec.getCode();
}
public BizException(int code, String message) {
super(message);
this.code = code;
}
public int getCode() { return code; }
}
package com.example.demo.common.exception;
import com.example.demo.common.ApiResult;
import com.example.demo.common.ErrorCode;
import lombok.extern.slf4j.Slf4j;
import org.springframework.web.bind.MethodArgumentNotValidException;
import org.springframework.web.bind.annotation.ExceptionHandler;
import org.springframework.web.bind.annotation.RestControllerAdvice;
@Slf4j
@RestControllerAdvice
public class GlobalExceptionHandler {
@ExceptionHandler(BizException.class)
public ApiResult<Void> handleBiz(BizException e) {
log.warn("业务异常:code={}, msg={}", e.getCode(), e.getMessage());
return ApiResult.fail(e.getCode(), e.getMessage());
}
@ExceptionHandler(MethodArgumentNotValidException.class)
public ApiResult<Void> handleValid(MethodArgumentNotValidException e) {
String msg = e.getBindingResult().getFieldErrors().stream()
.findFirst()
.map(f -> f.getField() + " " + f.getDefaultMessage())
.orElse(ErrorCode.PARAM_ERROR.getMessage());
return ApiResult.fail(ErrorCode.PARAM_ERROR.getCode(), msg);
}
@ExceptionHandler(Exception.class)
public ApiResult<Void> handleOther(Exception e) {
// 兜底异常对外只给模糊提示,详细堆栈只进日志,避免泄露内部细节
log.error("系统异常", e);
return ApiResult.fail(ErrorCode.SYSTEM_ERROR);
}
}
5.12 运行步骤
# 1. 启动中间件(见「三、环境准备」的 docker 命令)
# 2. 建库建表:执行 5.3 的 SQL + 第 6 章 6.3 的 local_message 建表 SQL
# 3. 启动应用
mvn spring-boot:run
# 4. 下单
curl -X POST http://localhost:8080/order/create
-H "Content-Type: application/json"
-d '{"userId":1001,"amount":99.90}'
# 5. 观察现象:
# - 控制台:MessageRelayTask 每 5 秒投递消息
# - 控制台:StockConsumer 手动 ACK 扣库存,NotifyService 用 async- 线程异步发通知
# - 30 分钟后:订单仍为 UNPAID 时,CloseConsumer 收到死信并关闭订单
验证清单(跑通即算掌握 L1/L2):
- 下单后,
orders表和local_message
表同时出现记录(同事务); local_message.status从 0 变 1(消息投递成功);- 扣库存日志出现一次,重复投递时出现「重复扣库存,已跳过」(幂等生效);
- 通知日志的线程名是
async-1(异步生效),扣库存日志是
SimpleAsyncTaskExecutor之外的自定义线程(手动 ACK
生效); - 30 分钟后,UNPAID 订单被关单(死信队列生效)。
6. 常见坑与排错指南
| 坑 / 现象 | 原因 | 解决方案 |
|---|---|---|
| 消息丢了(生产发出去没反应) | 生产者没开 confirm,发失败无感知 | 开启 publisher-confirm-type: correlated +publisher-returns,注册回调记日志告警 |
| 消费失败也把消息删了 | 用了默认自动 ACK(acknowledge-mode: auto) |
改 manual,成功 basicAck,失败basicNack 进死信 |
| 消费者突然收不到新消息 | 忘了 basicAck,消息一直「未确认」占满 prefetch额度 |
每个分支都确保 ack/nack,代码 review 重点检查 |
| 同一条消息被处理多次、数据错乱 | 消费端没做幂等 | 唯一键去重 / Redis setNX / DB 唯一索引 / 状态机 |
@Async 不生效,还是同步执行 |
同类自调用(this.xxx())没走代理 |
拆成两个 Bean,跨 Bean 调用 |
@Async 用了默认线程池导致 OOM |
没指定线程池,落到 SimpleAsyncTaskExecutor每任务新建线程 |
@Async("线程池名") + 自定义线程池 |
| 线程池任务堆积不处理 | 无界队列 + 无拒绝策略,任务无限排队 | 有界队列 + CallerRunsPolicy 拒绝策略 |
| 定时任务多实例重复执行 | 每个节点都跑同一个 @Scheduled |
Redis 分布式锁(setNX + Lua 释放)或 XXL-Job |
| 多个定时任务互相阻塞 | @Scheduled 默认单线程调度 |
配 TaskScheduler 线程池 |
| 死信队列延迟任务「不延迟」 | 延迟队列被挂了消费者,消息没等 TTL 就被取走 | 延迟队列不要绑定任何消费者,只监听死信队列 |
| 本地消息表消息不投递 | 消息表和业务表不在同一个数据库,事务罩不住 | 保证「同库同事务」,这是模式前提 |
8. 总结与延伸阅读
本篇围绕「消息与异步」这一主题,先讲透了 MQ 的价值(异步 / 解耦 /
削峰)和它引入的代价——三大可靠性问题(丢失 / 重复 /
顺序),并给出贯穿全篇的解决主线:丢→确认与持久化,重→幂等,序→单分区或业务规避。随后分别落地到
RabbitMQ(交换机路由 + confirm + 手动 ACK +
死信队列延迟任务)、Kafka(Partition / Consumer Group / Offset +
手动提交 + 幂等)、@Async(自定义线程池 +
失效场景)、定时任务(TaskScheduler 线程池 +
分布式锁)、分布式事务(本地消息表 + 事务消息 +
Seata)五条线,最后用一个「下单异步解耦」的完整实战把它们串起来。
记住三句话,你就能应对绝大多数 MQ
与异步相关的面试和事故:消息要确认与持久化才不丢,消费要幂等才不重,要保序就单分区或业务规避。
延伸阅读:
- RabbitMQ
官方文档 —— 交换机、确认机制、死信队列的权威定义; - Apache Kafka
官方文档 —— Partition、Consumer Group、Offset 的权威讲解; - 《深入理解 Kafka:核心设计与实践原理》(朱忠华)—— Kafka
底层原理的经典书籍; - XXL-Job 官方文档 ——
分布式任务调度的首选框架; - Seata
官方文档 —— AT / TCC / Saga 三种模式的原理与选型。
(阶段 6 完 · 下一篇将进入「阶段 7」)