【阶段 6】消息与异步:RabbitMQ、Kafka 与异步任务

86次阅读
没有评论

阶段 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. 学习目标与前置要求

学完你能…

  1. 说出消息队列的「三大可靠性问题」(丢失 / 重复 /
    顺序)分别如何解决,并能画出解决方案对照表;
  2. 独立写出 RabbitMQ 的生产者确认(confirm)、手动
    ACK、死信队列延迟任务三段完整代码;
  3. 独立写出 Kafka 手动提交 offset + 幂等消费的完整代码,并说清
    Partition / Consumer Group / Offset 三者关系;
  4. 独立配置一个自定义线程池并正确使用
    @Async,能验证并解释「同类自调用导致失效」的原因;
  5. 独立写出一个带 Redis 分布式锁的定时任务,避免多实例重复执行;
  6. 独立实现「本地消息表」模式,保证业务数据和消息数据同库同事务、消息不丢。

前置依赖

  • 阶段 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>

说明:ApiResultErrorCodeBizExceptionGlobalExceptionHandler
等公共类是贯穿全系列的「统一基础设施」,定义与前面阶段逐字一致(见阶段
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)
典型场景 业务消息、任务分发、延迟关单 日志采集、流处理、埋点、大数据量

选型口诀

  1. 业务消息(下单、扣款、发通知、延迟任务、复杂路由)→

    RabbitMQ。它路由灵活、延迟任务原生支持、运维简单,业务语义清晰。
  2. 日志 / 埋点 /
    流式数据
    (海量数据、需要回溯、需要把同一类数据喂给多个下游)→
    Kafka。它吞吐高、能回溯、天然支持多消费者组。
  3. 事务消息(本地事务与消息严格一致)→ 阿里系场景选
    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.createdorder.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 章的「丢」):

  1. 生产者确认(Publisher
    Confirm)
    :生产者发消息后,Broker
    回执告诉它「我收到了」,收不到回执就知道可能丢了;
  2. Broker
    持久化(durable)
    :交换机、队列、消息都持久化到磁盘,Broker
    重启不丢;
  3. 消费者手动 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 回调 + 路由失败退回

下面是一个生产级的订单消息生产者,它做了三件事:

  1. 通过 RabbitTemplate 发送 JSON 消息;
  2. 通过 ConfirmCallback 感知「Broker
    是否收到」,收不到则告警并补偿;
  3. 通过 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)?
一条消息变成死信有三种情况:

  1. 消息被消费者拒绝basicNack /
    basicRejectrequeue=false);
  2. 消息在队列里超时(超过 TTL,Time To Live
    存活时间);
  3. 队列长度满了,溢出的消息被丢弃。

什么是死信队列(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 关键点提示

  1. 手动 ACK 别漏 basicAck:忘了
    ack,消息会一直「未确认」占用 prefetch
    额度,最终消费者收不到新消息,队列堆积。
  2. 失败重试要设上限:配置里
    retry.max-attempts=3,重试 3 次仍失败再走
    basicNack
    进死信。否则一条坏消息会无限重试,拖垮消费端。
  3. 延迟队列千万别给它挂消费者:一旦有人消费
    order.delay.queue,消息没等到 TTL
    就被拿走了,延迟失效。
  4. 死信兜底不等于业务兜底:关单消费者里必须重新查库判断订单状态,因为「延迟的
    30 分钟里用户可能已支付」。
  5. 队列 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 关键点提示

  1. ack-modeenable-auto-commit
    要配套
    :想手动提交,就必须
    enable-auto-commit=falseack-mode=manual
    manual_immediate,否则自动提交会抢先提交 offset。
  2. 手动提交才有幂等问题的完整语境:手动提交 +
    处理失败不提交 → 消息重读 →
    重复,所以「手动提交」和「幂等」是一对必须同时出现的搭档。
  3. auto-offset-reset 的含义:只在「当前
    group 还没有提交过 offset」时生效,earliest
    从最早读,latest 从最新读。已有 offset 的 group
    不受它影响。
  4. 分区保序的前提:只有「同一 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;
    }
}

线程池的工作流程(面试必问,和参数强相关):

  1. 任务来了,先看核心线程(8
    个)有没有空闲:有空闲就交给核心线程执行;
  2. 核心线程都忙,任务进队列(容量 500)排队;
  3. 队列也满了,再看能否扩容到最大线程数(16
    个);
  4. 最大线程也忙、队列也满,触发拒绝策略

四种拒绝策略对比:

拒绝策略 行为 适用
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 关键点提示

  1. 线程池参数要按业务测算:核心线程数 ≈ CPU 核数(IO
    密集可放宽到 2×CPU 核数),队列容量按「高峰期任务积压量」估,拒绝策略用
    CallerRunsPolicy 兜底。
  2. 异步方法不要依赖 Request
    上下文
    @Async
    线程不是请求线程,RequestContextHolderThreadLocal
    里的登录态、traceId 都取不到,需要的话要显式传递(或使用
    TaskDecorator 透传)。
  3. 异步异常不会抛回调用方@Async
    方法的异常在独立线程抛出,调用方 catch 不到,所以异步方法内部必须自己
    try-catch + 记日志。
  4. 验证异步一定要看线程名:这是判断「到底异步没异步」的黄金标准,别靠感觉。

本章小结

@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 的 setNXsetIfAbsent),配合 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("报表生成完成");
    }
}

两个细节必须注意

  1. 锁要设过期时间(这里 30
    分钟),否则执行实例崩溃,锁永远不释放,任务永远被跳过;
  2. 释放锁要用 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 关键点提示

  1. fixedRate vs
    fixedDelay
    fixedRate
    是「固定频率」,任务执行时间超过间隔时会并发重叠fixedDelay
    是「固定间隔」,一定等上次结束才计时,不会重叠。要防重叠用
    fixedDelay
  2. 默认单线程是大坑:多个 @Scheduled
    必须配 TaskScheduler 线程池。
  3. 多实例必须加分布式锁:否则每个实例都执行一遍,扣款类任务会造成生产事故。
  4. 锁的过期时间要大于任务最坏执行时间,否则锁提前过期,任务可能被并发执行。

本章小结

@Scheduled 提供 cron / fixedRate / fixedDelay
三种触发方式,但默认单线程会导致任务互相阻塞,多实例部署会导致任务重复执行。生产上必须配调度线程池
+ 分布式锁(Redis setNX + Lua 释放),任务多了再上 XXL-Job
这类框架。


第 6 章
分布式事务(最终一致性)

6.1 为什么单库事务不够了

在阶段 4 里,你学会了用 @Transactional
保证「一个数据库里多个操作要么全成功、要么全失败」。这套机制在单体单库下完美,但一旦业务拆成多个服务、多个数据库,它就失效了。

看一个典型场景——下单:

订单服务(数据库 A):
    begin
    插入订单表  ✓
    发送「订单已创建」消息给 MQ
    commit

问题出在哪?「插入订单」和「发送消息」这两个操作无法放进同一个数据库事务里——数据库事务只能管数据库里的操作,管不了发消息。于是出现两种灾难:

  1. 先提交数据库,后发消息:如果 commit
    之后、发消息之前,进程崩溃了 → 订单写进去了,但消息没发出去 →
    库存服务永远不知道这笔订单,业务数据与消息不一致
  2. 先发消息,后提交数据库:如果发消息之后、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
内建版」,原理是「半消息 + 状态回查」:

  1. 生产者先发一条半消息(half
    message),此时消费者看不到它;
  2. 生产者执行本地事务(写订单);
  3. 本地事务成功 → 提交半消息,消费者可见;失败 →
    回滚半消息,消费者永远看不到;
  4. 如果生产者执行本地事务后没来得及提交/回滚就崩溃了,MQ
    定时回查生产者「你那个本地事务到底成没成」,据此决定提交还是回滚半消息。

它省去了本地消息表和扫描任务,但只限 RocketMQ
支持
,且需要处理回查逻辑。阿里系业务首选。

Seata
是一个分布式事务中间件,通过「全局事务协调器(TC)+
各服务分支事务(TM/RM)」实现跨服务一致性,有三种模式:

模式 原理 特点
AT 自动生成反向 SQL(undo log),自动回滚 无侵入,强一致,但性能损耗大
TCC 每个服务实现 Try/Confirm/Cancel 三个方法 手动编码,性能好,侵入大
Saga 长事务拆成多个本地事务,失败反向补偿 适合长流程,无锁

选型建议:本地消息表能解决 90% 的「业务 +
消息」一致性问题,先掌握它;Seata
只在「跨多个服务、多个数据库、要强一致」的复杂场景才引入,学习成本和运维成本都不低。

6.5 关键点提示

  1. 本地消息表必须「同库同事务」:如果消息表和业务表不在同一个数据库,本地事务就罩不住它们了,模式失效。这是理解这个模式的前提。
  2. 投递任务要有「状态机 +
    重试上限」
    status 是状态机(待投递 →
    已投递),retry_count
    用于告警和兜底,防止一条坏消息无限重试。
  3. 消费端幂等是最终闭环:本地消息表保证「不丢」,幂等保证「不重」,两者缺一不可。
  4. 消息体要存 JSONpayload 存序列化后的
    JSON 字符串,投递时反序列化,方便跨语言和排查。

本章小结

单库事务管不了「数据库 +
MQ」的跨资源一致性,业界用最终一致性解决。本地消息表是最经典、最通用的方案:业务数据和消息记录同库同事务写入,定时任务扫描投递,消费端幂等兜底。事务消息(RocketMQ)是它的内建版,Seata
用于跨服务强一致场景。


5.
生产级实战项目:下单异步解耦(完整可运行)

5.1 项目目标与架构

这个实战项目把本篇 80% 以上的知识点串成一条完整链路:

  1. 下单接口:本地事务里「写订单 +
    写本地消息表」,保证消息不丢(第 6 章);
  2. 消息投递任务:定时扫描消息表,按业务类型投递到
    RabbitMQ 的不同队列(第 6 章 + 第 2 章);
  3. 库存消费端:手动 ACK + Redis 幂等,扣减库存(第 2
    章);
  4. 通知消费端@Async 异步发通知邮件(第
    4 章);
  5. 延迟关单:TTL 30 分钟 +
    死信队列,未支付自动关单(第 2 章);
  6. 定时任务: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.queueTTL 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):

  1. 下单后,orders 表和 local_message
    同时出现记录(同事务);
  2. local_message.status 从 0 变 1(消息投递成功);
  3. 扣库存日志出现一次,重复投递时出现「重复扣库存,已跳过」(幂等生效);
  4. 通知日志的线程名是 async-1(异步生效),扣库存日志是
    SimpleAsyncTaskExecutor 之外的自定义线程(手动 ACK
    生效);
  5. 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
与异步相关的面试和事故:消息要确认与持久化才不丢,消费要幂等才不重,要保序就单分区或业务规避。

延伸阅读:

  1. RabbitMQ
    官方文档
    —— 交换机、确认机制、死信队列的权威定义;
  2. Apache Kafka
    官方文档
    —— Partition、Consumer Group、Offset 的权威讲解;
  3. 《深入理解 Kafka:核心设计与实践原理》(朱忠华)—— Kafka
    底层原理的经典书籍;
  4. XXL-Job 官方文档 ——
    分布式任务调度的首选框架;
  5. Seata
    官方文档
    —— AT / TCC / Saga 三种模式的原理与选型。

(阶段 6 完 · 下一篇将进入「阶段 7」)

正文完
 0
评论(没有评论)