阶段 11 · 从工程师走向架构师,解决分布式复杂问题 版本:Spring Boot
3.2.x / JDK 17 / Redis 7.x / MySQL 8.0
1. 导语
学到这里,你已经掌握了 Spring Boot
的绝大部分「单机」能力:分层架构、参数校验、MyBatis-Plus、事务、缓存、消息队列、微服务、测试、容器化部署。但当你把这些能力真正放到生产环境,把服务从
1 个实例扩到 10 个实例,把一张表从 100 万行写到 5000
万行,把一次下单从「前端点一下」变成「10 万人同时抢 1000
件商品」时,你会发现一批全新的问题接踵而至:
- 两台机器同时扣同一件商品的库存,单机的
synchronized
完全失效; - 用户手抖连点两次「提交订单」,系统扣了两次钱;
- 分库分表之后,原来数据库自增的主键 ID 开始重复;
- 一次下单要同时改库存、订单、积分、优惠券,跨服务后「事务」不再是你熟悉的那个事务;
- 业务越写越乱,Controller 里塞满了上千行的
if-else,改一处崩三处。
这些问题没有一个是「多学一个注解」能解决的,它们都属于分布式系统设计与代码架构的范畴,也正是「会写
CRUD 的工程师」与「能设计系统的架构师」之间的分水岭。
本篇是全系列的收尾篇。你将学会六块核心能力:分布式锁(含防死锁、防误删的完整推导)、幂等设计(接口幂等
+ 消息幂等 + 状态机)、分布式
ID(雪花算法结构与时钟回拨处理)、高并发架构(缓存/异步/限流/分库分表)、分布式事务一致性(CAP/BASE
与最终一致性落地)、DDD
领域驱动设计(限界上下文/聚合/值对象/领域事件,以及「什么时候该用、什么时候别用」的务实判断)。学完后,你会用一个「抢购」场景把这些能力串起来,并能独立设计一个分库分表的电商后端。
一句话预告:本篇没有新框架的「甜点」,全是架构判断力与工程决策——它们才是你未来几年薪资曲线的斜率。
2. 学习目标与前置要求
学完本篇,你能:
- 独立实现一个生产级分布式锁,并解释清楚「为什么加锁必须带超时」「为什么解锁必须用
Lua 校验 value」,能说出 Redisson 看门狗续期的原理,能对比
Redis、ZooKeeper、数据库三种锁方案。 - 为一个下单/扣款接口设计完整的幂等方案,覆盖唯一索引、状态机、请求号、Redis
setNX、乐观锁五种手段,并能说明「消息幂等」与「接口幂等」的差异。 - 手写一个雪花算法 ID 生成器,讲清 64
位每一位的含义,并能处理时钟回拨问题;能对比
UUID/雪花/号段/Redis INCR 四类方案的优劣。 - 为高并发系统设计缓存 + 异步 + 限流 + 读写分离 +
分库分表的组合方案,说清「什么时候该分库分表、按什么键怎么分」,并能说出分库分表后的三个连锁问题。 - 讲清 CAP/BASE
的区别,并能用本地消息表落地一次最终一致性的分布式事务;能对比本地消息表、TCC、Saga
三类方案的适用场景。 - 用 DDD
分层重构一个复杂订单模块,说清限界上下文、聚合、聚合根、值对象、领域事件的定义,并给出「简单
CRUD 别硬上 DDD」的明确判断标准。
前置依赖:
- 阶段 4(数据访问:MyBatis-Plus、事务、缓存)——本篇大量复用事务与
Redis 基础; - 阶段
6(消息队列:RabbitMQ、死信队列、消息可靠性)——本篇幂等与最终一致性会串联消息幂等; - 阶段
7(微服务:Nacos/Feign/Gateway/Sentinel)——本篇限流与分布式事务在微服务语境下展开。
若你还没掌握上述阶段,建议先回看对应文章;尤其是阶段 4 的
@Transactional 传播行为与阶段 6
的消息可靠投递,本篇会直接引用、不再重复推导。
3. 环境准备
本篇示例依赖 Redis 与 MySQL。请确认环境版本如下:
| 组件 | 版本 | 说明 |
|---|---|---|
| JDK | 17 | 全系列统一 |
| Spring Boot | 3.2.x | 全系列统一 |
| Redis | 7.x | 分布式锁、幂等、分布式 ID 的基础 |
| MySQL | 8.0 | 唯一索引、乐观锁、本地消息表 |
| Redisson | 3.27.2 | 生产级分布式锁首选 |
| ShardingSphere-JDBC | 5.4.x | 分库分表(第 4 章概念 + 配置) |
| RabbitMQ | 3.12.x | 消息幂等、最终一致性(复用阶段 6) |
pom.xml 关键依赖(在阶段 4 的基础上增量添加):
<dependencies>
<!-- 阶段 4 已有的:web / validation / mybatis-plus / mysql / lombok -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<!-- Redisson:生产级分布式锁,内置看门狗续期 -->
<dependency>
<groupId>org.redisson</groupId>
<artifactId>redisson-spring-boot-starter</artifactId>
<version>3.27.2</version>
</dependency>
<!-- AOP:幂等注解切面需要 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-aop</artifactId>
</dependency>
<!-- ShardingSphere-JDBC:分库分表(第 4 章演示,可选引入) -->
<dependency>
<groupId>org.apache.shardingsphere</groupId>
<artifactId>shardingsphere-jdbc-core</artifactId>
<version>5.4.1</version>
</dependency>
</dependencies>
application.yml 关键配置:
spring:
data:
redis:
host: localhost
port: 6379
datasource:
url: jdbc:mysql://localhost:3306/dk_mall?useSSL=false&serverTimezone=Asia/Shanghai
username: root
password: root
初始化命令(若 Redis 未启动):
# 启动 Redis(Docker)
docker run -d --name redis -p 6379:6379 redis:7
# 启动 MySQL(Docker,密码 root)
docker run -d --name mysql -p 3306:3306 -e MYSQL_ROOT_PASSWORD=root mysql:8.0
本篇复用阶段 3/4
的公共基础设施(ApiResult、ErrorCode、BizException、GlobalExceptionHandler),不再重新定义。其中
ErrorCode
在本篇按需扩展了三个错误码,基础码数值保持不变:
public enum ErrorCode {
SUCCESS(0, "success"),
PARAM_ERROR(40001, "参数错误"),
UNAUTHORIZED(40101, "未登录或登录已过期"),
FORBIDDEN(40301, "无权限访问"),
USER_NOT_FOUND(40401, "用户不存在"),
DUPLICATE_SUBMIT(40901, "重复提交,请勿重复操作"),
TOO_MANY_REQUESTS(42901, "请求过于频繁,请稍后重试"),
SOLD_OUT(41001, "商品已售罄"),
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; }
}
说明:
DUPLICATE_SUBMIT(40901)、TOO_MANY_REQUESTS(42901)、SOLD_OUT(41001)是本篇新增的扩展码,数值不与基础码冲突。
4. 正文章节
第 1 章:分布式锁
1.1 为什么需要分布式锁
先回顾一个你早已熟悉的场景。单机下,要防止两个线程同时扣库存,你用
synchronized 或 ReentrantLock 就够了:
// 单机互斥:JVM 内的锁,只能锁住同一进程内的线程
public class SingleMachineStock {
private final ReentrantLock lock = new ReentrantLock();
private int stock = 100;
public void deduct() {
lock.lock();
try {
if (stock <= 0) {
throw new BizException(ErrorCode.SOLD_OUT);
}
stock--;
} finally {
lock.unlock();
}
}
}
这段代码在单实例下没有任何问题。但一旦你的服务为了扛住流量部署了
3 个实例(负载均衡后面挂 3
台机器),问题就来了:ReentrantLock 是 JVM 进程内的锁,实例
A 的锁根本管不到实例 B 的线程。3 个实例各扣各的,stock
还是会被扣成负数——这就是经典的超卖问题。
分布式锁要解决的问题:多个进程(多台机器、多个
JVM)之间,对同一个共享资源(一段缓存、一张表的某一行、一个文件)实现互斥访问。
一个合格的分布式锁,必须同时满足四个条件:
- 互斥性:同一时刻只能有一个客户端持有锁(这是锁的底线)。
- 不会死锁:即使持有锁的进程崩溃、没来得及释放,锁也必须能自动过期,不能永久卡死。
- 不会误删:锁只能由「持有它的那个客户端」释放,不能出现「A
的锁被 B 释放」。 - 高可用:锁服务(Redis)挂了要能切到备机,不能因为锁服务故障让整个业务不可用。
第 1 和第 2 条对应「防死锁」,第 3
条对应「防误删」,这两点是本章的重点,也是面试和线上事故的高发区。本章的讲解顺序,就是一把分布式锁从「能用」到「好用」的完整演进路线:
SET NX EX 占坑 → Lua 校验解锁(防误删) → 看门狗续期(防锁过期) → Redisson(生产封装)
1.2 最小可用:Redis SET NX EX
Redis
之所以适合做分布式锁,核心在于它的单线程命令模型:每一个命令执行期间不会被其他命令打断。基于此,我们可以用一条原子命令完成「占坑」:
SET lock:order:1001 <唯一值> NX EX 30
NX(Not eXists):只有当 key 不存在时才写入,返回
OK;key 已存在则返回 nil。这就是「占坑」的原子判断。EX 30:设置 30
秒过期时间。这是「防死锁」的关键——万一持锁进程崩溃,30
秒后锁自动释放,不会永久卡死。
NX 和 EX
必须放在同一条命令里,这是有讲究的。如果拆成两条命令(先
SETNX 再
EXPIRE),一旦在两步之间进程崩溃,就会留下一个「永不过期」的锁,直接死锁。这正是总规范里反复强调的「关键陷阱」,注释里必须讲清为什么。
用 Spring Boot 的 StringRedisTemplate
实现最小可用版本:
@Slf4j
@Service
@RequiredArgsConstructor
public class SimpleDistributedLock {
private final StringRedisTemplate redisTemplate;
/**
* 尝试加锁(最小可用版)。
*
* @param lockKey 锁的 key,一般是「业务前缀 + 业务 id」,如 lock:order:1001
* @param requestId 本次请求的唯一标识,解锁时校验「是不是我加的锁」用
* @param expireSeconds 锁的超时时间,防止持锁进程崩溃后死锁
* @return true=加锁成功,false=别人正持有锁
*/
public boolean tryLock(String lockKey, String requestId, long expireSeconds) {
// setIfAbsent 对应 SET NX EX:NX 保证互斥,EX 保证防死锁,两者原子完成
Boolean ok = redisTemplate.opsForValue()
.setIfAbsent(lockKey, requestId, expireSeconds, TimeUnit.SECONDS);
return Boolean.TRUE.equals(ok);
}
}
requestId 用 UUID
生成,每个线程/每次请求一个,作为「这把锁是谁加的」的凭证:
String requestId = UUID.randomUUID().toString();
这个版本已经能用,但它有一个致命缺陷:没有正确解锁。请看
1.3。
1.3 误删问题与 Lua 解锁
先看一个「错误解锁」会引发的线上事故,这是本章必须讲透的误删问题。
假设锁的超时时间是 30 秒,但业务逻辑执行了 40 秒(GC 停顿、慢
SQL、下游抖动都会导致):
- 线程 A 拿到锁,开始执行,业务跑了 40 秒;
- 到第 30 秒,A 的锁过期自动释放了;
- 线程 B 在第 31 秒成功拿到锁,开始执行;
- 第 40 秒,A 的业务执行完了,调用
DEL lockKey
想「释放自己的锁」; - 结果 A 把 B 正在持有的锁删掉了;
- 线程 C 立刻拿到了锁——此时 B 和 C 同时进入临界区,互斥被打破。
这就是「误删」:A 删掉了 B
的锁。根因是「判断锁是不是我的」和「删除锁」这两个动作不是原子的。如果解锁代码写成这样:
// 反例:先 GET 判断,再 DEL 删除,两步之间可能被插队,导致误删
public void unlockWrong(String lockKey, String requestId) {
String value = redisTemplate.opsForValue().get(lockKey);
if (requestId.equals(value)) { // 判断:锁是我的
redisTemplate.delete(lockKey); // 删除:但这两步之间,锁可能已经过期并被别人拿走
}
}
GET 和 DEL 之间,锁可能过期、可能被 B
拿走,A 依然会把 B 的锁删掉。要解决,就必须让「判断 +
删除」变成一个原子操作。Redis
原生命令做不到「带判断的删除」,但 Lua 脚本可以——Lua
脚本在 Redis 服务端执行,执行期间不会被其他命令穿插,天然原子。
解锁 Lua 脚本:
-- KEYS[1] = 锁的 key,ARGV[1] = 当前请求的 requestId
-- 只有当锁的 value 等于我的 requestId 时才删除,否则返回 0(不是我加的锁,不动它)
if redis.call('get', KEYS[1]) == ARGV[1] then
return redis.call('del', KEYS[1])
else
return 0
end
Java 侧用 DefaultRedisScript 执行:
@Slf4j
@Service
@RequiredArgsConstructor
public class RedisDistributedLock {
private final StringRedisTemplate redisTemplate;
/** 解锁 Lua 脚本:校验 value + 删除,两步在 Redis 服务端原子完成,防止误删别人的锁 */
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
);
public boolean tryLock(String lockKey, String requestId, long expireSeconds) {
Boolean ok = redisTemplate.opsForValue()
.setIfAbsent(lockKey, requestId, expireSeconds, TimeUnit.SECONDS);
return Boolean.TRUE.equals(ok);
}
/**
* 解锁:用 Lua 保证「校验 requestId + 删除」原子。
* 返回 true 表示成功释放;返回 false 表示锁已过期或被别人持有(不处理,避免误删)。
*/
public boolean unlock(String lockKey, String requestId) {
Long result = redisTemplate.execute(UNLOCK_SCRIPT,
Collections.singletonList(lockKey), requestId);
return result != null && result == 1L;
}
}
到这里,你手上的这个锁已经解决了「互斥、防死锁、防误删」三个核心问题,可以应付很多场景了。但它还剩最后一个隐患:锁的过期时间定多长?
- 定太短(比如 3
秒):业务还没执行完锁就过期了,别的线程进来,互斥失效; - 定太长(比如 30 分钟):持锁进程一旦崩溃,这 30
分钟内所有人都得等,业务被「卡死」半分钟。
这个矛盾靠「人肉估一个超时时间」是无法根治的,必须让锁「自动续期」——这就是
1.4 要讲的看门狗。
1.4 看门狗续期与 Redisson
看门狗(Watchdog)
的思路是:持锁期间,由一个后台线程定时给锁「续命」;只要业务还在执行,锁就永远不会过期;业务正常结束,锁被主动释放;进程崩溃,后台线程也没了,锁在最后一次续期后自然过期。这样既不会因为业务慢而误删,也不会因为进程崩溃而死锁。
这个机制自己写很繁琐(要维护一个后台续期线程池、要处理续期失败、要处理可重入计数),生产上直接使用
Redisson——它是 Java 生态里最成熟的 Redis
客户端之一,把「分布式锁 + 看门狗 + 可重入 +
自动续期」全部封装好了。
先配置 Redisson(单机模式,生产可用哨兵/集群模式):
@Configuration
public class RedissonConfig {
@Bean
public RedissonClient redissonClient() {
Config config = new Config();
// 单机模式;生产建议替换为 useSentinelServers() 或 useClusterServers() 以支持高可用
config.useSingleServer()
.setAddress("redis://localhost:6379")
.setConnectionPoolSize(16);
return Redisson.create(config);
}
}
Redisson 加锁的标准姿势:
@Slf4j
@Service
@RequiredArgsConstructor
public class RedissonLockService {
private final RedissonClient redissonClient;
public void doWithLock(String lockKey, Runnable business) {
RLock lock = redissonClient.getLock(lockKey);
boolean locked = false;
try {
// tryLock(waitTime, leaseTime, unit):
// waitTime = 3s:最多等 3 秒拿不到锁就放弃,避免请求无限堆积
// leaseTime = -1(不传第三个参数,默认 -1):启用看门狗,默认 30 秒租期,每 10 秒自动续期
locked = lock.tryLock(3, TimeUnit.SECONDS);
if (!locked) {
// 拿不到锁说明并发竞争激烈,直接快速失败,让上层走限流/降级
throw new BizException(ErrorCode.TOO_MANY_REQUESTS);
}
// 拿到锁,执行业务
business.run();
} catch (InterruptedException e) {
// 等待锁期间被打断,恢复中断标记并抛业务异常
Thread.currentThread().interrupt();
throw new BizException(ErrorCode.SYSTEM_ERROR);
} finally {
// isHeldByCurrentThread():只释放「当前线程持有的锁」,防止误删其他线程/其他实例的锁
if (locked && lock.isHeldByCurrentThread()) {
lock.unlock();
}
}
}
}
关于 Redisson 的几个关键点,必须理解清楚:
- 看门狗只在
leaseTime == -1
时生效。tryLock(3, TimeUnit.SECONDS)只传了
waitTime,leaseTime 走默认值 -1,此时 Redisson 会给锁一个默认 30
秒的租期,并启动看门狗,每 10 秒(租期的 1/3)自动续期一次,直到主动
unlock。 - 如果显式指定了 leaseTime(如
tryLock(3, 30, TimeUnit.SECONDS)),则不启用看门狗,锁
30 秒后必过期——适合你明确知道业务一定会在 30 秒内完成的场景。 isHeldByCurrentThread()的判断很关键:Redisson
的锁是可重入的,如果 A 线程加了两次锁,第一次的unlock
只是把重入计数减一,并没有真正释放。用这个判断可以避免「还没真正释放就提前走
finally」的误操作。- 可重入:Redisson 内部用 Hash 结构记录「线程 id →
重入次数」,同一线程可以重复加锁。
Redisson 还提供了多种锁类型,按场景选用:
| 锁类型 | 获取方式 | 适用场景 |
|---|---|---|
| 普通锁 | getLock(key) |
绝大多数互斥场景 |
| 公平锁 | getFairLock(key) |
要求按请求先后顺序获取锁 |
| 读写锁 | getReadWriteLock(key) |
读多写少,读读不互斥、读写互斥 |
| 联锁 | getMultiLock(lock1, lock2) |
需要同时锁多个资源 |
1.5 生产级完整示例:扣减库存
把上面的知识串成一个生产级的扣库存方法——加锁、执行业务、finally
释放、拿不到锁快速失败:
@Slf4j
@Service
@RequiredArgsConstructor
public class StockService {
private final RedissonClient redissonClient;
private final StockMapper stockMapper;
/**
* 扣减库存(分布式锁保护)。
* 场景:多实例部署下,防止两个实例同时扣同一商品库存导致超卖。
*/
public void deductStock(Long productId, int quantity) {
String lockKey = "lock:stock:" + productId;
RLock lock = redissonClient.getLock(lockKey);
boolean locked = false;
try {
// 等 2 秒,拿不到锁说明该商品正在被其他请求处理,快速失败
locked = lock.tryLock(2, TimeUnit.SECONDS);
if (!locked) {
log.warn("获取分布式锁失败, productId={}", productId);
throw new BizException(ErrorCode.TOO_MANY_REQUESTS);
}
// 临界区:检查库存 + 扣减。锁内仍需乐观锁兜底(见第 2 章幂等与第 4 章)
int stock = stockMapper.selectStock(productId);
if (stock < quantity) {
throw new BizException(ErrorCode.SOLD_OUT);
}
stockMapper.deduct(productId, quantity);
log.info("扣减库存成功, productId={}, quantity={}", productId, quantity);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new BizException(ErrorCode.SYSTEM_ERROR);
} finally {
// 只释放当前线程持有的锁,防止误删
if (locked && lock.isHeldByCurrentThread()) {
lock.unlock();
}
}
}
}
关键提示:分布式锁只能保证「临界区在同一时刻只有一个实例在跑」,它不能替代数据库层面的并发控制。锁
+
乐观锁(update ... where num >= quantity)双层防护,才是生产上的标准答案,这也是第
5 章实战项目会组合使用的原因。
1.5.1
锁粒度与性能的三个实战细节
分布式锁用起来简单,但「锁得对不对」直接决定系统的吞吐和正确性。三个实战细节值得单独展开:
(1)锁的粒度要尽可能小。锁
lock:stock:1001(锁单个商品)而不是
lock:stock(锁全部商品),因为前者让「抢不同商品的请求」互不阻塞,后者会让所有抢购请求串行化,吞吐直接崩掉。锁的
key
设计原则:锁住「会发生冲突的最小共享资源」,而不是一大片资源。
(2)锁内代码要尽可能短。锁持有时间越长,其他请求等待越久。锁内只放「必须互斥的临界区」(读库存
+
扣库存),耗时的动作(调外部接口、发消息、写日志到磁盘)一律移到锁外:
// 锁内只做临界区,耗时动作移出锁外
public void deductStockOptimized(Long productId, int quantity) {
RLock lock = redissonClient.getLock("lock:stock:" + productId);
boolean locked = false;
try {
locked = lock.tryLock(2, TimeUnit.SECONDS);
if (!locked) {
throw new BizException(ErrorCode.TOO_MANY_REQUESTS);
}
// 临界区(短):一条乐观锁 SQL 搞定,不查库存再扣(省一次往返)
int rows = stockMapper.deduct(productId, quantity);
if (rows <= 0) {
throw new BizException(ErrorCode.SOLD_OUT);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new BizException(ErrorCode.SYSTEM_ERROR);
} finally {
if (locked && lock.isHeldByCurrentThread()) {
lock.unlock();
}
}
// 锁外(耗时):发通知、写审计日志、发消息,不占用锁
notifyService.sendDeductNotice(productId, quantity);
}
(3)拿不到锁要有「降级路径」。tryLock
拿不到锁时,不要无限重试(会让请求堆积、线程占满),而是快速失败 +
上层限流/降级兜底,或者排队一段时间后放弃。tryLock(waitTime, ...)
里的 waitTime 就是「最多等多久」,超过就返回
false,这是避免雪崩的关键设计。
1.6
三种分布式锁方案对比与选型
除了 Redis,业界还有 ZooKeeper 和数据库两种实现,各有取舍:
| 方案 | 原理 | 优点 | 缺点 |
|---|---|---|---|
| Redis(Redisson) | SET NX EX + 看门狗 | 性能高、生态成熟、上手快 | 极端网络分区下可能「锁丢失」(需 Redlock 缓解) |
| ZooKeeper | 临时顺序节点 + Watch | 强一致、锁释放靠会话失效,天然防死锁 | 性能较低、依赖 ZK 集群、引入额外组件 |
| 数据库 | SELECT ... FOR UPDATE / 唯一索引 |
简单、不引入新组件 | 性能最差、需手动处理过期 |
选型建议:绝大多数业务用 Redis + Redisson
即可(性能与易用性最佳);对一致性要求极高、可接受性能损耗的场景(如资金、配置中心)用
ZooKeeper。关于 Redis 分布式锁在「主从切换丢锁」的极端情况,Redlock
算法(向多个独立 Redis
实例加锁,多数成功才算成功)提供了一种缓解方案,但也存在争议——生产上通常用「单实例
Redis + 哨兵高可用」配合业务幂等兜底,性价比更高。
本章小结:分布式锁的演进路线是「SET NX EX 占坑 → Lua
校验解锁防误删 →
看门狗续期防锁过期」,每前进一步都是在解决前一步暴露的真实事故。Redisson
是生产首选,但你必须理解它背后就是这三步,否则出了问题只会盲目重启。
第 2 章:幂等设计
2.1 为什么必须幂等
「幂等」这个词源于数学(f(f(x)) = f(x)),在软件里的含义是:同一个操作执行一次和执行多次,产生的业务结果完全相同。
为什么在分布式系统里幂等是「必选项」而不是「加分项」?因为分布式环境下的「重试」无处不在:
- 网络超时重试:下单请求发出去,服务端其实已经扣款成功了,但响应在网络里丢了,客户端超时后自动重试,导致重复扣款。
- 用户重复点击:手抖连点两下「提交」,或前端按钮没做防抖,同一笔订单提交了两次。
- 消息重复投递:RabbitMQ 的
at-least-once
投递语义下,消费者宕机后消息会重新投递,你无法保证每条消息只消费一次。 - 失败重试框架:你主动引入的重试机制(Feign
重试、RocketMQ 消费重试、Job 调度重试)。
结论是:在分布式系统里,你永远无法保证「请求只来一次」,你能保证的只有「请求来多次结果也不变」。所以幂等不是「要不要做」,而是「在哪个层面做、用什么方案做」。
一个典型的非幂等接口——直接插入订单:
// 反例:非幂等。同一请求重试一次,就多插入一条订单
public void createOrder(OrderDTO dto) {
Order order = new Order();
order.setOrderNo(dto.getOrderNo());
order.setAmount(dto.getAmount());
orderMapper.insert(order); // 每次调用都插一条,重试就重复
}
2.2 五种幂等方案
| 方案 | 实现方式 | 适用场景 | 特点 |
|---|---|---|---|
| 唯一索引 | 数据库唯一约束(如 order_no 唯一) |
天然有业务唯一键的场景 | 最可靠,但依赖 DB,且唯一键必须选对 |
| 状态机 | 校验状态流转是否合法(如「已支付」不能再「待支付」) | 订单、审批等有明确状态的业务 | 把幂等藏在业务规则里 |
| 请求号(Token) | 前端先取 token,后端校验后删除 | 表单防重复提交 | 需要前后端配合,一次性 |
| Redis setNX | 先 setIfAbsent 占坑,占不到就拒绝 |
接口级幂等,通用性强 | 简单灵活,但依赖 Redis |
| 乐观锁 | update ... where version = ? 或where num > 0 |
并发更新同一行 | 天然防并发,适合扣库存/改状态 |
下面逐个给出生产级实现。
2.3 唯一索引 + 状态机
唯一索引是最简单也最可靠的幂等手段。以订单为例,只要
order_no(订单号)有唯一约束,重复插入同一订单号就一定会被数据库拒绝:
CREATE TABLE t_order (
id BIGINT PRIMARY KEY,
order_no VARCHAR(64) NOT NULL,
user_id BIGINT NOT NULL,
amount DECIMAL(10,2) NOT NULL,
status TINYINT NOT NULL DEFAULT 0,
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
UNIQUE KEY uk_order_no (order_no) -- 唯一索引:幂等的最后一道防线
) ENGINE=InnoDB;
@Slf4j
@Service
@RequiredArgsConstructor
public class OrderService {
private final OrderMapper orderMapper;
public void createOrder(OrderDTO dto) {
try {
orderMapper.insert(toEntity(dto));
} catch (DuplicateKeyException e) {
// 唯一索引冲突 = 该订单已存在,说明是重复请求,直接吞掉当作成功(幂等)
log.info("订单已存在,忽略重复提交, orderNo={}", dto.getOrderNo());
}
}
}
状态机适合订单这类「状态有向图」的业务。把「当前状态
→
下一状态」的合法流转固化下来,非法流转直接拒绝,重复的合法流转也自然幂等:
public enum OrderStatus {
PENDING(0), // 待支付
PAID(1), // 已支付
SHIPPED(2), // 已发货
CLOSED(3); // 已关闭
private final int code;
OrderStatus(int code) { this.code = code; }
// 状态机:定义每个状态允许流转到的目标状态
private static final Map<OrderStatus, Set<OrderStatus>> ALLOWED = Map.of(
PENDING, Set.of(PAID, CLOSED),
PAID, Set.of(SHIPPED),
SHIPPED, Set.of(),
CLOSED, Set.of()
);
public boolean canTransitTo(OrderStatus target) {
return ALLOWED.getOrDefault(this, Set.of()).contains(target);
}
public static OrderStatus of(int code) {
return Arrays.stream(values()).filter(s -> s.code == code)
.findFirst().orElseThrow(() -> new BizException(ErrorCode.PARAM_ERROR));
}
}
public void markPaid(Long orderId) {
Order order = orderMapper.selectById(orderId);
OrderStatus current = OrderStatus.of(order.getStatus());
// 状态机校验:已支付的订单不能再次「已支付」,天然幂等
if (!current.canTransitTo(OrderStatus.PAID)) {
log.info("订单状态不允许流转,忽略重复操作, orderId={}, current={}", orderId, current);
return;
}
// 用乐观锁更新,where status = 旧状态,防止并发下两个线程同时流转成功
orderMapper.updateStatus(orderId, current.getCode(), OrderStatus.PAID.getCode());
}
对应的状态机 SQL(where status = 旧状态 本身就是乐观锁 +
幂等的组合):
@Update("UPDATE t_order SET status = #{target} WHERE id = #{id} AND status = #{current}")
int updateStatus(@Param("id") Long id, @Param("current") int current, @Param("target") int target);
2.4 请求号幂等(Redis setNX)
唯一索引和状态机都需要「业务天然有唯一键」,但很多接口没有(比如「发送验证码」「提交评价」)。这时用请求号
+ Redis setNX 做接口级幂等最通用:调用方每次请求带一个唯一的
requestNo,服务端先占坑,占得到才执行。
先定义入参 DTO(带 requestNo):
@Data
public class SubmitOrderDTO {
/** 请求唯一标识:客户端生成(如 UUID),同一笔业务多次重试必须传同一个值 */
@NotBlank(message = "requestNo 不能为空")
private String requestNo;
@NotNull(message = "商品 id 不能为空")
private Long productId;
@NotNull(message = "数量不能为空")
@Min(value = 1, message = "数量至少为 1")
private Integer quantity;
}
幂等的核心实现:
@Slf4j
@Service
@RequiredArgsConstructor
public class IdempotentSubmitService {
private final StringRedisTemplate redisTemplate;
private static final String IDEMPOTENT_KEY_PREFIX = "idempotent:order:";
/** 幂等记录保留 10 分钟:覆盖「客户端重试」的合理时间窗口,避免 key 永久堆积 */
private static final long EXPIRE_MINUTES = 10;
public ApiResult<Void> submit(SubmitOrderDTO dto) {
String key = IDEMPOTENT_KEY_PREFIX + dto.getRequestNo();
// setIfAbsent 原子占坑:只有第一个请求能占成功,后续同名 requestNo 全部被拒
Boolean ok = redisTemplate.opsForValue()
.setIfAbsent(key, "1", EXPIRE_MINUTES, TimeUnit.MINUTES);
if (!Boolean.TRUE.equals(ok)) {
log.info("重复请求被拦截, requestNo={}", dto.getRequestNo());
throw new BizException(ErrorCode.DUPLICATE_SUBMIT);
}
try {
doSubmit(dto); // 真正执行业务
return ApiResult.ok(null);
} catch (Exception e) {
// 业务失败必须删掉占坑的 key,否则用户修正参数后重试会被误判为「重复提交」
redisTemplate.delete(key);
throw e;
}
}
private void doSubmit(SubmitOrderDTO dto) {
// 实际下单逻辑:扣库存、插订单……
log.info("下单成功, requestNo={}, productId={}", dto.getRequestNo(), dto.getProductId());
}
}
这段代码有两个必须理解的设计决策(也是面试高频考点):
- 为什么失败要删
key:如果不删,业务因为参数错误/库存不足失败后,用户改好参数再次提交,会被误判成「重复提交」永远提交不了。删掉
key 表示「这次尝试作废,允许重试」。 - 为什么占坑成功但业务执行到一半进程崩溃:此时 key
没删、业务也没做成,会「卡住」这个 requestNo 最长 10
分钟。这个时间窗口的取舍是工程权衡——太短挡不住慢重试,太长会误伤。对支付类高价值业务,通常配合「唯一索引」兜底:即使
requestNo 机制漏了,数据库唯一约束也会拦住重复。
2.5 乐观锁幂等
乐观锁本质上是「把幂等和并发控制合并成一条
SQL」。它不依赖锁、不依赖唯一索引,靠 version
字段或业务条件(如
num > 0)来保证「同一行只能被成功更新一次」:
@Mapper
public interface StockMapper {
/** 乐观锁扣库存:where stock >= quantity 保证「库存够才扣」。返回 0 表示失败,天然防超卖 */
@Update("UPDATE t_stock SET stock = stock - #{quantity}, version = version + 1 " +
"WHERE product_id = #{productId} AND stock >= #{quantity}")
int deduct(@Param("productId") Long productId, @Param("quantity") int quantity);
}
乐观锁适合「更新同一行」的场景(扣库存、改订单状态),它的优点是不引入
Redis、性能好(一条 SQL 搞定);缺点是需要业务表有可判断的字段(version
或业务条件),且不适合「插入」类操作。
2.6 注解 + AOP 实现通用幂等
上面的 setNX 幂等如果每个接口都手写一遍「占坑 + 失败删
key」,会很啰嗦。生产上常封装成注解 + AOP
切面,业务代码只加一个 @Idempotent 注解即可。
定义注解:
@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
public @interface Idempotent {
/** 幂等 key 前缀,如 "order" */
String prefix();
/** 幂等记录过期时间,单位秒 */
long expireSeconds() default 600;
}
切面实现(通过 SpEL 或参数名取 requestNo):
@Slf4j
@Aspect
@Component
@RequiredArgsConstructor
public class IdempotentAspect {
private final StringRedisTemplate redisTemplate;
@Around("@annotation(idempotent)")
public Object around(ProceedingJoinPoint joinPoint, Idempotent idempotent) throws Throwable {
// 从方法参数中提取 requestNo(约定 DTO 里必须有 requestNo 字段)
String requestNo = extractRequestNo(joinPoint.getArgs());
if (requestNo == null) {
throw new BizException(ErrorCode.PARAM_ERROR.getCode(), "幂等接口缺少 requestNo");
}
String key = "idempotent:" + idempotent.prefix() + ":" + requestNo;
Boolean first = redisTemplate.opsForValue()
.setIfAbsent(key, "1", idempotent.expireSeconds(), TimeUnit.SECONDS);
if (!Boolean.TRUE.equals(first)) {
throw new BizException(ErrorCode.DUPLICATE_SUBMIT);
}
try {
return joinPoint.proceed(); // 真正执行被注解的方法
} catch (Throwable e) {
redisTemplate.delete(key); // 业务失败删 key,允许重试
throw e;
}
}
private String extractRequestNo(Object[] args) {
for (Object arg : args) {
if (arg == null) {
continue;
}
try {
Method getter = arg.getClass().getMethod("getRequestNo");
Object value = getter.invoke(arg);
if (value != null) {
return value.toString();
}
} catch (NoSuchMethodException ignored) {
// 不是带 requestNo 的 DTO,跳过
} catch (Exception e) {
throw new BizException(ErrorCode.SYSTEM_ERROR);
}
}
return null;
}
}
业务代码只需加注解,幂等逻辑完全透明:
@Idempotent(prefix = "order", expireSeconds = 600)
public ApiResult<Void> submit(SubmitOrderDTO dto) {
// 业务逻辑,无需关心幂等
return ApiResult.ok(null);
}
2.7 消息幂等
消息队列的幂等是另一条主线(串起阶段 6)。RabbitMQ
的消费者处理到一半宕机,消息会重新投递,消费者会收到同一条消息多次。因此消费端必须幂等。
方案:给每条消息一个全局唯一的
messageId(生产端生成,或消费端用业务唯一键),消费端用「唯一索引
+ 先查后插」或「Redis setNX」去重:
@Slf4j
@Component
@RequiredArgsConstructor
public class OrderMessageListener {
private final OrderMapper orderMapper;
private final StringRedisTemplate redisTemplate;
@RabbitListener(queues = "order.paid.queue")
public void onMessage(OrderPaidMessage msg) {
String dedupKey = "msg:dedup:" + msg.getMessageId();
// 消费端幂等:用 messageId 占坑,重复投递直接丢弃
Boolean first = redisTemplate.opsForValue()
.setIfAbsent(dedupKey, "1", 24, TimeUnit.HOURS);
if (!Boolean.TRUE.equals(first)) {
log.info("重复消息被丢弃, messageId={}", msg.getMessageId());
return;
}
try {
processPaidOrder(msg);
} catch (Exception e) {
// 处理失败删掉去重 key,让消息重试时能重新处理
redisTemplate.delete(dedupKey);
throw e;
}
}
private void processPaidOrder(OrderPaidMessage msg) {
log.info("处理订单支付消息, orderNo={}", msg.getOrderNo());
// 实际业务:把订单状态改为已支付(配合状态机幂等)
}
}
接口幂等 vs
消息幂等的差异,是自测题的考点:接口幂等面对的是「同一请求的重复提交」,靠请求号去重;消息幂等面对的是「同一消息的重复投递」,靠消息
id 去重。二者手段相似(都是唯一标识 +
占坑/唯一索引),但触发来源不同,且消息幂等还需要和「消息确认(ACK)」配合——处理成功才
ACK,失败不 ACK 让 MQ 重投,重投又靠幂等兜底。
本章小结:幂等的本质是「把『请求只能来一次』的不可能,转化为『请求来多次结果不变』的可能」。唯一索引和状态机藏在数据层,最可靠;请求号和
setNX
藏在入口层,最通用;乐观锁藏在更新层,专治并发。生产上通常是多层叠加,而不是二选一。
第 3 章:分布式 ID
3.1 为什么自增 ID 不够用了
单库单表时代,主键用数据库自增(AUTO_INCREMENT)是天经地义的:简单、有序、还自带「趋势递增」的索引友好特性。但一旦你开始分库分表(第
4 章),问题立刻暴露:
- 订单表按
user_id哈希拆成了 4 张表,每张表都有独立的
AUTO_INCREMENT,从 1 开始自增。于是user_id=3
和user_id=7的两个订单,主键都是
1001,主键在全局重复了。 - 而主键是很多业务关联(订单明细表、退款表)的外键,主键一重复,关联关系就乱了。
所以分布式环境下,主键必须全局唯一。一个合格的分布式
ID 要满足:
- 全局唯一:任何两张表、任何两个实例生成的 ID
都不能重复。 - 趋势递增:对 InnoDB
这类聚簇索引,主键有序能显著提升插入性能(避免页分裂),也方便排序和范围查询。 - 高性能:高并发下生成 ID
不能成为瓶颈(要能到每秒几十万甚至百万级)。 - 不含敏感信息:自增 ID
会暴露业务量(竞争对手看到「今日第 9999999
单」就知道了你的规模),分布式 ID 通常要避免这种规律性。
3.2 四种方案对比
| 方案 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| UUID | 本地生成,零依赖,绝对唯一 | 无规律字符串、太长(36 字符/128 位)、无序导致索引页分裂严重 | 非主键的临时标识、会话 id |
| 雪花算法(Snowflake) | 趋势递增、高性能、本地生成无网络开销 | 依赖机器时钟,时钟回拨需处理 | 生产主键首选 |
| 数据库号段 | 简单可控、可批量预取 | 依赖 DB,存在单点(可用双号段缓解) | 中小规模、已有 DBA 体系 |
| Redis INCR | 简单、天然递增 | 依赖 Redis,持久化丢失会重复 | 对 ID 连续性有要求的场景 |
为什么雪花算法比 UUID
更适合做主键,是自测题的固定考点,答案核心在两点:UUID
无序,写入 InnoDB
时随机插入导致频繁页分裂,索引性能随数据量增长急剧恶化;雪花算法按时间递增,天然贴合聚簇索引的插入顺序,且长度只有
64 位(long),比 36 字符的字符串 UUID 小得多,存储和索引都更省。
3.3 雪花算法结构
雪花算法把 64 位 long 拆成四段:
┌─┬──────────────────────────┬──────────────┬──────────────────┐
│1│ 41 bit │ 10 bit │ 12 bit │
│符│ 时间戳(毫秒) │ 机器 ID │ 序列号 │
│号│ 相对自定义起始时间偏移 │ (数据中心+机器)│ 同一毫秒内自增 │
└─┴──────────────────────────┴──────────────┴──────────────────┘
63 62 22 12 0
- 1 bit 符号位:恒为 0,保证生成的 ID 是正数(long
最高位为 1 是负数)。 - 41 bit 时间戳:记录「当前毫秒 –
自定义起始时间」的偏移量。41 位能表示约 69.7
年(2^41 / 1000 / 3600 / 24 / 365 ≈ 69.7),从 2020
年起算可以用到 2089 年。 - 10 bit 机器 ID:标识是哪台机器生成的,最多 1024
台机器。通常再拆成「5 bit 数据中心 + 5 bit 机器」。 - 12 bit 序列号:同一毫秒内自增,支持每毫秒生成 4096
个 ID;超出则等待下一毫秒。
算一下上限:单机每毫秒 4096 个,即每秒约 409 万个
ID,足以扛住绝大多数业务的峰值。
用一个具体例子理解拼接过程。假设 workerId=1,起始时间戳为
EPOCH,当前毫秒偏移为 ts,序列号为
seq,那么:
// 时间戳左移 22 位,占据 62~22 位;机器 ID 左移 12 位,占据 21~12 位;序列号占 11~0 位
long id = (ts << 22) | (workerId << 12) | seq;
三段「按位或」拼到一起,得到一个 64 位、趋势递增、全局唯一的
long。
3.4 最小可用实现
public class SimpleSnowflakeIdGenerator {
/** 起始时间戳:2020-01-01 00:00:00,自定义,用于压缩时间戳位数 */
private static final long EPOCH = 1577836800000L;
private static final long WORKER_ID_BITS = 10L;
private static final long SEQUENCE_BITS = 12L;
private static final long MAX_WORKER_ID = ~(-1L << WORKER_ID_BITS); // 1023
private static final long SEQUENCE_MASK = ~(-1L << SEQUENCE_BITS); // 4095
private static final long WORKER_ID_SHIFT = SEQUENCE_BITS; // 12
private static final long TIMESTAMP_SHIFT = SEQUENCE_BITS + WORKER_ID_BITS; // 22
private final long workerId;
private long lastTimestamp = -1L;
private long sequence = 0L;
public SimpleSnowflakeIdGenerator(long workerId) {
if (workerId < 0 || workerId > MAX_WORKER_ID) {
throw new IllegalArgumentException("workerId 必须在 0~1023 之间");
}
this.workerId = workerId;
}
public synchronized long nextId() {
long timestamp = System.currentTimeMillis();
if (timestamp < lastTimestamp) {
// 时钟回拨:当前时间比上次生成时间还小,说明系统时钟被往回拨了
throw new IllegalStateException("系统时钟回拨,拒绝生成 ID");
}
if (timestamp == lastTimestamp) {
// 同一毫秒:序列号自增,超出 4095 则自旋等待下一毫秒
sequence = (sequence + 1) & SEQUENCE_MASK;
if (sequence == 0) {
timestamp = waitNextMillis(lastTimestamp);
}
} else {
sequence = 0L; // 新的毫秒,序列号归零
}
lastTimestamp = timestamp;
// 拼接:时间戳左移 22 位 | 机器 ID 左移 12 位 | 序列号
return ((timestamp - EPOCH) << TIMESTAMP_SHIFT)
| (workerId << WORKER_ID_SHIFT)
| sequence;
}
/** 自旋等待,直到时钟追上 lastTimestamp */
private long waitNextMillis(long lastTimestamp) {
long timestamp = System.currentTimeMillis();
while (timestamp <= lastTimestamp) {
timestamp = System.currentTimeMillis();
}
return timestamp;
}
}
3.5 时钟回拨问题(生产级)
上面的最小实现遇到时钟回拨会直接抛异常,这在生产上是不可接受的——时钟回拨虽然罕见,但确实会发生:NTP
校时、运维手动改时间、虚拟机迁移,都可能让系统时钟往回跳几毫秒甚至几秒。一旦回拨,雪花算法就可能生成重复
ID。
生产级处理时钟回拨有几种策略:
- 短时间回拨等待:回拨在几毫秒内,自旋等待时钟追上来(最小实现已含此逻辑,但只针对「序列号溢出」场景,要扩展成「回拨也等待」)。
- 超过阈值拒绝:回拨超过一定毫秒数(如
5ms),说明不是抖动而是真回拨,此时抛异常或降级。 - 备用 workerId:检测到回拨后,临时切换到一个备用机器
ID 生成,等时钟恢复再切回。 - 用外部单调时钟兜底:如从 Redis
取当前时间,保证单调递增。
下面是增强版:回拨小于 5ms 时自旋等待追平,超过 5ms
直接拒绝并打日志告警:
@Slf4j
public class RobustSnowflakeIdGenerator {
private static final long EPOCH = 1577836800000L;
private static final long WORKER_ID_BITS = 10L;
private static final long SEQUENCE_BITS = 12L;
private static final long SEQUENCE_MASK = ~(-1L << SEQUENCE_BITS);
private static final long WORKER_ID_SHIFT = SEQUENCE_BITS;
private static final long TIMESTAMP_SHIFT = SEQUENCE_BITS + WORKER_ID_BITS;
/** 允许的最大回拨毫秒数,超过则判定为异常回拨 */
private static final long MAX_BACKWARD_MS = 5L;
private final long workerId;
private long lastTimestamp = -1L;
private long sequence = 0L;
public RobustSnowflakeIdGenerator(long workerId) {
if (workerId < 0 || workerId > ~(-1L << WORKER_ID_BITS)) {
throw new IllegalArgumentException("workerId 必须在 0~1023 之间");
}
this.workerId = workerId;
}
public synchronized long nextId() {
long timestamp = currentTimeMillis();
// 时钟回拨处理:小于 5ms 的抖动,自旋等待追平;否则拒绝并告警
if (timestamp < lastTimestamp) {
long offset = lastTimestamp - timestamp;
if (offset <= MAX_BACKWARD_MS) {
log.warn("检测到时钟回拨 {}ms,自旋等待追平", offset);
timestamp = waitUntil(lastTimestamp);
} else {
log.error("时钟回拨超过阈值 {}ms,拒绝生成 ID,请检查 NTP/系统时钟", offset);
throw new IllegalStateException("时钟回拨异常,拒绝生成 ID");
}
}
if (timestamp == lastTimestamp) {
sequence = (sequence + 1) & SEQUENCE_MASK;
if (sequence == 0) {
timestamp = waitUntil(lastTimestamp + 1);
}
} else {
sequence = 0L;
}
lastTimestamp = timestamp;
return ((timestamp - EPOCH) << TIMESTAMP_SHIFT)
| (workerId << WORKER_ID_SHIFT)
| sequence;
}
private long waitUntil(long target) {
long now = currentTimeMillis();
while (now < target) {
now = currentTimeMillis();
}
return now;
}
/** 单独抽一个方法便于测试时 mock 时间 */
protected long currentTimeMillis() {
return System.currentTimeMillis();
}
}
3.6 号段模式(数据库发号)
如果你不想处理时钟回拨,又希望 ID
有序,号段模式是另一个生产可用的选择:数据库维护一个「当前最大
ID」记录,每次用 UPDATE ... SET max_id = max_id + step
原子地「租」一段 ID(如一次 1000
个),应用在内存里逐个发放,发完再租下一段。
号段表结构:
CREATE TABLE t_id_segment (
biz VARCHAR(32) PRIMARY KEY, -- 业务标识,如 "order"
max_id BIGINT NOT NULL, -- 当前已分配到的最大 ID
step INT NOT NULL DEFAULT 1000 -- 每次租的号段大小
) ENGINE=InnoDB;
号段生成器(内存发号 + 双号段预取):
@Slf4j
@Component
@RequiredArgsConstructor
public class SegmentIdGenerator {
private final IdSegmentMapper segmentMapper;
/** 当前正在发放的号段 [start, end] 与当前位置 */
private volatile long currentStart = 0;
private volatile long currentEnd = 0;
private final AtomicLong cursor = new AtomicLong(0);
public synchronized long nextId(String biz) {
// 当前号段用尽,去数据库原子地租下一段
if (cursor.get() >= currentEnd) {
loadSegment(biz);
}
return cursor.incrementAndGet();
}
/** 原子地租一段号:UPDATE 让并发线程不会租到重叠的区间 */
private void loadSegment(String biz) {
// 原子取号段:UPDATE ... SET max_id = max_id + step WHERE biz = ?
// 数据库返回更新后的 max_id,则本段区间为 (max_id - step, max_id]
long newMax = segmentMapper.getAndIncrement(biz);
long step = segmentMapper.getStep(biz);
this.currentStart = newMax - step + 1;
this.currentEnd = newMax;
this.cursor.set(currentStart - 1);
log.info("号段已加载, biz={}, range=[{}, {}]", biz, currentStart, currentEnd);
}
}
@Mapper
public interface IdSegmentMapper {
/** 原子取号段:同一条 SQL 完成「读 + 加 step」,并发下不会租到重复区间 */
@Update("UPDATE t_id_segment SET max_id = max_id + step WHERE biz = #{biz}")
int getAndIncrement(@Param("biz") String biz);
@Select("SELECT step FROM t_id_segment WHERE biz = #{biz}")
int getStep(@Param("biz") String biz);
}
号段模式的优点是「趋势递增、无时钟依赖、可批量预取性能高」;缺点是依赖数据库(有单点,可用双号段「当前段
+ 后台线程预取下一段」缓解,避免发号线程阻塞在数据库上)。美团 Leaf
等中间件就是号段 + 雪花两种模式都支持的代表。
3.6.1 Redis INCR 方案
如果你要的 ID 不需要「趋势递增」那么讲究,只是「全局唯一 +
单调递增」,Redis 的 INCR 是最简单的方案:每次调用原子自增
1:
@Slf4j
@Component
@RequiredArgsConstructor
public class RedisIdGenerator {
private final StringRedisTemplate redisTemplate;
/** Redis INCR 原子自增,天然全局唯一 + 单调递增 */
public long nextId(String biz) {
Long id = redisTemplate.opsForValue().increment("id:seq:" + biz);
return id == null ? 0L : id;
}
}
Redis INCR 的致命弱点是持久化丢失:如果 Redis 用 RDB
快照且没开 AOF,重启后会丢失一部分已分配的 ID,导致重复;即使开了
AOF(每秒刷盘),宕机也可能丢失最后一秒的增量。所以它适合「对 ID
连续性无要求、可接受极低概率重复」的场景(如临时编号、消息流水号),不适合做主键。真要用于主键,通常给
INCR 加一个「起始偏移量」且开启 AOF 并容忍边界丢失。
3.7 生产落地:MyBatis-Plus
内置雪花
自己写雪花算法是为了理解原理(面试、排查 ID
重复问题时用得上),生产落地不必重复造轮子——MyBatis-Plus
已经内置了雪花算法,配合主键策略即可:
@Data
@TableName("t_order")
public class Order {
// ASSIGN_ID:MyBatis-Plus 内置雪花算法,插入时自动生成趋势递增的全局唯一 ID
@TableId(type = IdType.ASSIGN_ID)
private Long id;
private String orderNo;
private Long userId;
}
MyBatis-Plus 的默认实现同样是 64 位雪花(时间戳 + workerId +
序列号),workerId 默认根据 MAC 地址自动分配,也可通过
mybatis-plus.global-config.db-config.id-type 或自定义
IdentifierGenerator 配置。
本章小结:分布式 ID
的核心矛盾是「全局唯一」与「有序、高性能」的平衡。雪花算法用 64
位把时间、机器、序列号编码进一个
long,同时满足三者,代价是引入「时钟回拨」这个必须处理的边界情况。理解它的结构,比会调
API 重要得多。
第 4 章:高并发架构
4.1 高并发系统的核心手段全景
「高并发」不是某一个技术,而是一组手段的组合。面对流量,架构师手里的牌大致五张,按「离用户由近到远」排列:
- 缓存:读多写少的场景,用 Redis 挡住 90%
的读流量,减轻数据库压力(Cache Aside 模式)。 - 异步:写操作走消息队列异步化,把「同步等待」变成「削峰填谷」,让瞬时流量被队列缓冲。
- 限流/降级/熔断:在入口处拒绝超额请求(限流),核心链路挂了返回兜底(降级),下游故障快速失败(熔断),保护系统不被冲垮。
- 读写分离:主库写、从库读,把读流量分散到多个从库。
- 分库分表:当单库单表成为物理瓶颈(数据量、连接数、写入量)时,水平拆分。
前四张牌你其实在前面的阶段已经见过:缓存是阶段 4 的 Redis,异步是阶段
6 的 RabbitMQ,限流熔断是阶段 7 的 Sentinel,读写分离是阶段 4
的多数据源。本章把重点放在限流的实现细节和分库分表这两块最「架构」的部分,其余做全景串讲。
4.2 缓存:Cache Aside
模式与三大经典问题
缓存最经典的使用方式是 Cache
Aside(旁路缓存):读时先查缓存,命中直接返回;未命中查库,回填缓存。写时先更新数据库,再删除缓存(而不是更新缓存)。
@Slf4j
@Service
@RequiredArgsConstructor
public class ProductCacheService {
private final ProductMapper productMapper;
private final StringRedisTemplate redisTemplate;
private static final String CACHE_KEY = "product:";
public Product getProduct(Long id) {
// 1. 先查缓存
String json = redisTemplate.opsForValue().get(CACHE_KEY + id);
if (json != null) {
return JSON.parseObject(json, Product.class);
}
// 2. 缓存未命中,查库
Product product = productMapper.selectById(id);
if (product == null) {
return null;
}
// 3. 回填缓存(带过期时间,防止 key 永久堆积)
redisTemplate.opsForValue().set(CACHE_KEY + id, JSON.toJSONString(product), 30, TimeUnit.MINUTES);
return product;
}
public void updateProduct(Product product) {
// 先更新数据库
productMapper.updateById(product);
// 再删除缓存(而不是更新缓存):下次读时重新加载,避免双写不一致
redisTemplate.delete(CACHE_KEY + product.getId());
}
}
关键点:「先改库、后删缓存」的顺序不能反。若「先删缓存再改库」,在删缓存和改库之间来了一个读请求,会把旧数据重新读进缓存,导致缓存里是脏数据。这是阶段
4 讲过的缓存一致性核心,此处不再展开。
缓存还有三个「经典问题」,高并发下必须会防(阶段 4
已详讲,这里只列对策):
| 问题 | 现象 | 对策 |
|---|---|---|
| 缓存穿透 | 查一个根本不存在的 key,每次都打穿缓存到 DB | 空值缓存 / 布隆过滤器 |
| 缓存击穿 | 热点 key 过期瞬间,大量请求同时打到 DB | 互斥锁重建 / 热点 key 永不过期 |
| 缓存雪崩 | 大量 key 同时过期,DB 瞬间被打崩 | 过期时间加随机值 / 多级缓存 |
三个问题的生产级防护代码:
@Slf4j
@Service
@RequiredArgsConstructor
public class CacheProtectedService {
private final StringRedisTemplate redisTemplate;
private final ProductMapper productMapper;
private static final String CACHE_KEY = "product:";
private static final String LOCK_KEY = "lock:rebuild:";
public Product getProduct(Long id) {
String key = CACHE_KEY + id;
String json = redisTemplate.opsForValue().get(key);
if (json != null) {
return "NULL".equals(json) ? null : JSON.parseObject(json, Product.class);
}
// 缓存未命中,查库并回填
Product product = productMapper.selectById(id);
if (product == null) {
// 防穿透:缓存空值("NULL"),下次查询直接命中缓存,不打 DB
redisTemplate.opsForValue().set(key, "NULL", 1, TimeUnit.MINUTES);
return null;
}
// 防雪崩:过期时间加随机值,避免大量 key 同时过期
long ttl = 30 + ThreadLocalRandom.current().nextInt(10);
redisTemplate.opsForValue().set(key, JSON.toJSONString(product), ttl, TimeUnit.MINUTES);
return product;
}
/** 防击穿:热点 key 过期时,用分布式锁保证只有一个线程重建缓存,其余等待 */
public Product getProductWithLock(Long id) {
String key = CACHE_KEY + id;
String json = redisTemplate.opsForValue().get(key);
if (json != null) {
return JSON.parseObject(json, Product.class);
}
// 只有一个线程能拿到重建锁,其他线程自旋重试读缓存
String lockKey = LOCK_KEY + id;
Boolean locked = redisTemplate.opsForValue().setIfAbsent(lockKey, "1", 5, TimeUnit.SECONDS);
if (Boolean.TRUE.equals(locked)) {
try {
// 双重检查:拿到锁后再读一次缓存,防止重复查库
json = redisTemplate.opsForValue().get(key);
if (json != null) {
return JSON.parseObject(json, Product.class);
}
Product product = productMapper.selectById(id);
redisTemplate.opsForValue().set(key, JSON.toJSONString(product), 30, TimeUnit.MINUTES);
return product;
} finally {
redisTemplate.delete(lockKey); // 重建完释放锁
}
} else {
// 没拿到锁:短暂等待后重试读缓存(生产可用重试几次)
return retryGetFromCache(key);
}
}
private Product retryGetFromCache(String key) {
for (int i = 0; i < 3; i++) {
String json = redisTemplate.opsForValue().get(key);
if (json != null) {
return JSON.parseObject(json, Product.class);
}
try { Thread.sleep(50); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
}
return null;
}
}
关键点:防击穿的重建锁和本篇第 1
章的分布式锁是同一套思想——用「只有一个能重建」避免缓存击穿瞬间 DB
被打穿。区别在于这里锁的是「缓存重建」这个动作,锁的粒度是单个 key。
4.3 异步削峰:消息队列
同步接口的吞吐上限是「数据库能扛多少」,异步化后,接口只需「把消息丢进队列就返回」,瞬时压力被
MQ 缓冲,消费者按数据库能承受的速度慢慢消费:
@Slf4j
@Service
@RequiredArgsConstructor
public class OrderAsyncService {
private final RabbitTemplate rabbitTemplate;
/** 下单异步化:接口只负责落库 + 发消息,耗时动作(发短信、加积分、通知仓库)全部异步 */
public void createOrderAsync(OrderDTO dto) {
// 1. 快速落库(这是同步链路里唯一必须同步的部分)
// orderMapper.insert(...);
// 2. 发消息,把「下单成功」这个事件广播出去,下游各自消费
rabbitTemplate.convertAndSend("order.exchange", "order.created", dto);
log.info("下单消息已发送, orderNo={}", dto.getOrderNo());
}
}
关于异步化,阶段 6 已经讲透了可靠投递(confirm/return
回调、死信队列关单),本篇只强调它在高并发架构里的定位:异步是削峰的工具,不是「让接口变快」的魔法——它把「不可控的瞬时压力」变成了「可控的积压」,代价是引入最终一致性(第
5 章)。
4.4 限流:固定窗口与令牌桶
限流的目标很朴素:宁可拒绝一部分请求,也不能让系统被压垮。常用的算法有计数器、滑动窗口、令牌桶、漏桶。先看最简单的固定窗口计数器——用
Redis + Lua 实现一个「每 key 每秒钟最多 N 个请求」的限流器:
@Slf4j
@Component
@RequiredArgsConstructor
public class RedisRateLimiter {
private final StringRedisTemplate redisTemplate;
/**
* 固定窗口限流 Lua 脚本:
* KEYS[1] = 限流 key,ARGV[1] = 窗口内允许的最大请求数,ARGV[2] = 窗口秒数
* 思路:key 不存在则创建并计数 1;存在且未超限则自增;超限则返回 0 拒绝。
*/
private static final DefaultRedisScript<Long> RATE_LIMIT_SCRIPT = new DefaultRedisScript<>(
"local current = redis.call('incr', KEYS[1]) " +
"if current == 1 then " +
" redis.call('expire', KEYS[1], ARGV[2]) " +
"end " +
"if current > tonumber(ARGV[1]) then " +
" return 0 " +
"end " +
"return 1",
Long.class
);
/** 尝试获取令牌,true=放行,false=被限流 */
public boolean tryAcquire(String key, int limit, int windowSeconds) {
Long result = redisTemplate.execute(RATE_LIMIT_SCRIPT,
Collections.singletonList(key), String.valueOf(limit), String.valueOf(windowSeconds));
return result != null && result == 1L;
}
}
固定窗口有个缺陷:窗口边界处会出现「双倍突发」(前一个窗口的最后一秒和后一个窗口的第一秒各满额,实际
2 秒内放行了 2N
个请求)。要更平滑,用滑动窗口或令牌桶。令牌桶允许一定突发、又限制平均速率,是生产上最常用的方案。单机场景可直接用
Guava 的 RateLimiter:
@Slf4j
@Component
public class GuavaRateLimiter {
/** 每秒允许 100 个请求(令牌桶:固定速率放令牌,允许突发取走桶内令牌) */
private final RateLimiter rateLimiter = RateLimiter.create(100.0);
public boolean tryAcquire() {
// tryAcquire:有令牌立即放行;没令牌立即返回 false(不阻塞等待)
return rateLimiter.tryAcquire();
}
}
分布式场景下,令牌桶需要用 Redis + Lua 实现(记录「上次补充令牌的时间
+ 当前令牌数」),或用阶段 7 讲过的 Sentinel
统一限流。在抢购接口入口处使用限流:
@RestController
@RequiredArgsConstructor
public class SeckillController {
private final RedisRateLimiter rateLimiter;
private final SeckillService seckillService;
@PostMapping("/seckill")
public ApiResult<Void> seckill(@RequestParam Long userId, @RequestParam Long productId) {
// 接口级限流:同一用户每秒最多 5 次,把刷子挡在业务逻辑外面
if (!rateLimiter.tryAcquire("rate:seckill:" + userId, 5, 1)) {
throw new BizException(ErrorCode.TOO_MANY_REQUESTS);
}
seckillService.seckill(userId, productId);
return ApiResult.ok(null);
}
}
4.5 读写分离
读写分离的本质是「把读和写拆到不同的数据库实例」:主库承担写(INSERT/UPDATE/DELETE),从库承担读(SELECT),主从之间通过
binlog 异步同步。
- 收益:读流量被分散到多个从库,主库专注写入,整体吞吐成倍提升。
- 代价:主从同步有延迟,刚写入的数据立刻读从库可能读不到(读己之写问题)。解决:关键读走主库、或按业务容忍度接受短暂不一致。
MyBatis-Plus 配合动态数据源(如
dynamic-datasource-spring-boot-starter)可实现按注解路由读写:
@Service
@RequiredArgsConstructor
public class OrderQueryService {
private final OrderMapper orderMapper;
// @DS("slave"):该查询路由到从库(需要引入 dynamic-datasource)
public Order getById(Long id) {
return orderMapper.selectById(id);
}
}
4.6
分库分表:什么时候分、怎么分
分库分表是五张牌里「最重」的一张,因为一旦拆分,SQL、事务、跨表查询、运维全都变复杂,能不拆就不拆,拆了就要一次拆对。判断标准有两条硬指标:
- 数据量:单表数据量超过 2000
万行(或索引都放不进内存、慢查询频发)时考虑分表; - 写入/连接瓶颈:单库连接数(MySQL
默认几百)被打满、写入 QPS 到上限时考虑分库。
怎么分,核心是选对分片键(Sharding
Key):
- 按用户 ID
哈希:库 = user_id % 4,表 = user_id % 8,即「4
库 8
表」。同一个用户的所有数据落在同一张表,用户维度的查询天然不跨表,这是电商用户中心最常用的分法。 - 按时间分表:日志、流水按月份分表,历史数据方便归档。
- 按订单号哈希:订单表按
order_no
哈希,适合订单维度的查询。
分片键的选择原则只有一条:让绝大多数查询都命中同一个分片,避免跨库
JOIN 和全分片扫描。
以用户表分「4 库 8 表」为例,ShardingSphere-JDBC 的配置:
spring:
shardingsphere:
datasource:
names: ds0, ds1, ds2, ds3
ds0:
type: com.zaxxer.hikari.HikariDataSource
driver-class-name: com.mysql.cj.jdbc.Driver
jdbc-url: jdbc:mysql://localhost:3306/dk_mall_0
username: root
password: root
# ds1 / ds2 / ds3 同理,指向 dk_mall_1 / 2 / 3
rules:
sharding:
tables:
t_user:
# 分库 + 分表:库用 user_id % 4,表用 user_id % 8
actual-data-nodes: ds$->{0..3}.t_user_$->{0..7}
database-strategy:
standard:
sharding-column: user_id
sharding-algorithm-name: db-inline
table-strategy:
standard:
sharding-column: user_id
sharding-algorithm-name: tbl-inline
sharding-algorithms:
db-inline:
type: INLINE
props:
algorithm-expression: ds$->{user_id % 4}
tbl-inline:
type: INLINE
props:
algorithm-expression: t_user_$->{user_id % 8}
props:
sql-show: true
分库分表后有三个「连锁反应」必须提前规划:
- 全局 ID:主键不能用自增了,改用第 3
章的雪花算法。 - 数据迁移:存量数据要按分片规则搬到新库新表(用双写
+ 迁移工具 + 灰度切流),这是分库分表最容易翻车的环节。 - 查询改造:跨片查询(如按手机号查用户,但分片键是
user_id)需要引入基因法(把 user_id
的某几位编码进手机号或订单号,反解出分片键)或 ES 等异构索引辅助。
基因法示例(解决「按非分片键查询」):订单号里嵌入
user_id % 8
的余数作为「基因」。这样即使前端只拿到订单号,也能从订单号末尾的基因反推出
user_id 所在的表,避免全分片扫描:
// 生成订单号时嵌入分片基因:orderNo 末位 = user_id % 8
public String buildOrderNo(Long userId) {
long gene = userId % 8;
// 订单号 = 时间戳 + 雪花序列 + 基因位,末位编码了 user_id 的取模结果
return System.currentTimeMillis() + String.valueOf(gene);
}
除了 YAML,ShardingSphere 也支持纯 Java
配置(便于按环境动态调整分片规则):
@Configuration
public class ShardingConfig {
@Bean
public DataSource dataSource() throws SQLException {
// 1. 数据源映射:4 个库
Map<String, DataSource> dataSourceMap = new HashMap<>();
for (int i = 0; i < 4; i++) {
dataSourceMap.put("ds" + i, createDataSource("jdbc:mysql://localhost:3306/dk_mall_" + i));
}
// 2. 分片规则:t_user 分库(user_id % 4)+ 分表(user_id % 8)
ShardingRuleConfiguration ruleConfig = new ShardingRuleConfiguration();
ruleConfig.getTables().add(getUserTableRule());
// 3. 全局主键策略:雪花算法(对应第 3 章)
ruleConfig.setKeyGenerators(Collections.singletonMap(
"snowflake", new AlgorithmConfiguration("SNOWFLAKE", new Properties())));
Properties props = new Properties();
props.setProperty("sql-show", "true");
return ShardingSphereDataSourceFactory.createDataSource(dataSourceMap,
Collections.singletonList(ruleConfig), props);
}
private ShardingTableRuleConfiguration getUserTableRule() {
ShardingTableRuleConfiguration table = new ShardingTableRuleConfiguration("t_user",
"ds$->{0..3}.t_user_$->{0..7}");
// 分库策略:按 user_id 哈希取模
table.setDatabaseShardingStrategy(new StandardShardingStrategyConfiguration(
"user_id", "db-inline"));
// 分表策略:按 user_id 哈希取模
table.setTableShardingStrategy(new StandardShardingStrategyConfiguration(
"user_id", "tbl-inline"));
// 雪花主键
table.setKeyGenerateStrategy(new KeyGenerateStrategyConfiguration("id", "snowflake"));
return table;
}
}
4.7
数据迁移:分库分表最危险的一步
分库分表的技术本身不难,难的是在不影响线上业务的前提下,把存量数据搬到新库新表。标准迁移流程分四步:
- 双写:改造写路径,让新数据同时写旧表和新分片(通过
Canal 等 binlog 同步工具,或应用层双写)。此时以旧表为准。 - 历史迁移:用离线任务把存量数据按分片规则批量搬到新库新表,跑完后做数据校验(行数、抽样对账)。
- 灰度切流:先切 1%
流量读新库新表,观察无误后逐步放量到 100%。此时以新分片为准。 - 下线旧表:稳定运行一段时间(如一周)后,停止双写,下线旧表。
整个过程中「什么时候以哪边为准」是核心——切流前以旧表为准(新表只是影子),切流后以新分片为准(旧表只读)。这也是为什么「分库分表要提前规划」——如果业务已经跑起来了才想起来分,迁移成本会成倍上升。
本章小结:高并发架构没有银弹,只有「根据瓶颈在哪、选对应的牌」。缓存治读多,异步治瞬时峰值,限流治刷子,读写分离治读压力,分库分表治数据量和写入天花板。五张牌按需组合,但每张都有代价——架构的本质是权衡。
第 5 章:分布式事务一致性
5.1 CAP 与
BASE:从「强一致」到「最终一致」
单机事务(ACID)的前提是「一个数据库、一个连接」。微服务化之后,一次下单要同时改「订单库、库存库、积分库、优惠券库」,跨了多个数据库甚至多个服务,ACID
的「原子性」物理上就不成立了。于是分布式系统必须重新回答「一致性」这个问题。
CAP
定理:一个分布式系统在「一致性(Consistency)、可用性(Availability)、分区容错性(Partition
tolerance)」三者中,最多只能同时满足两个。
- 一致性(C):所有节点同一时刻看到的数据都一样(强一致)。
- 可用性(A):每个请求都能在合理时间内得到响应。
- 分区容错性(P):网络分区(节点之间通信中断)时系统仍能工作。
关键洞察是:网络分区是分布式系统无法回避的客观现实,P
必须满足。于是现实的选择只剩两个——CP(牺牲可用性保一致,如
ZooKeeper、Etcd)或 AP(牺牲一致性保可用,如 Eureka、多数
NoSQL)。对电商这种「可用性优先」的业务,几乎必然选择
AP,也就是接受最终一致性。
BASE 理论是 AP 的具体化:
- Basically
Available(基本可用):故障时允许损失部分可用性(如降级、排队),但核心功能仍可用。 - Soft
state(软状态):允许系统存在中间状态,数据不必时刻一致。 - Eventually
consistent(最终一致性):不要求实时一致,但保证最终会一致(可能几秒后)。
一句话总结:CAP 告诉你「为什么不能既要又要」,BASE
告诉你「放弃强一致后怎么保证最终一致」。
用一个具体例子理解「最终一致」和「强一致」的差别。假设「用户支付成功后」,需要同时完成三件事:订单变已支付(订单库)、库存减一(库存库)、积分加
100(积分库):
- 强一致(2PC):三个库必须「要么都成功、要么都失败」,任何一个库失败,另外两个也回滚。代价是三个库全程锁住、同步阻塞,吞吐极低。
- 最终一致(BASE):先保证「订单变已支付」这个核心动作成功,然后通过消息把「已支付」这个事实异步广播给库存库和积分库,它们各自消费、各自更新。中间可能有几秒库存库还没更新,但最终一定会更新。代价是「短时间内订单和库存不一致」。
电商场景下,用户对「积分晚几秒到账」几乎无感,但对「支付卡住 30
秒」极度敏感,所以最终一致是电商的必然选择。理解了「代价是什么、用户能不能接受」,你才算真正懂了
CAP/BASE,而不是背定义。
5.2 四种分布式事务方案对比
| 方案 | 原理 | 一致性 | 复杂度 | 适用场景 |
|---|---|---|---|---|
| 2PC / XA | 两阶段提交,所有参与者 prepare 后统一 commit | 强一致 | 低(协议本身简单) | 跨库但同类型数据库、强一致必须 |
| 本地消息表 | 业务 + 消息表同库同事务,再异步投递 | 最终一致 | 中 | 电商下单、转账等通用场景 |
| TCC | Try(预留)→ Confirm(确认)→ Cancel(回滚) | 最终一致 | 高(每个接口写三遍) | 资金、支付等需精确控制的场景 |
| Saga | 长事务拆成一串本地事务 + 补偿动作 | 最终一致 | 中 | 长流程、跨多个服务编排 |
四者的取舍本质是:强一致(2PC)性能差、可用性低,最终一致(本地消息表/TCC/Saga)性能好但要接受短暂不一致并写好补偿。生产上,绝大多数业务选最终一致,其中本地消息表是性价比最高、最通用的入门方案。
5.3 本地消息表(生产级实现)
本地消息表(又称 Transactional
Outbox,事务发件箱)的核心思路:把「业务操作」和「要发送的消息」放进同一个本地事务里,保证二者要么都成功、要么都失败;事务提交后,再用一个后台任务把消息可靠地投递出去。
下单成功 + 记录「订单已创建」消息,两件事绑定:
@Slf4j
@Service
@RequiredArgsConstructor
public class OrderTxService {
private final OrderMapper orderMapper;
private final OutboxMessageMapper outboxMapper;
/** 下单 + 写消息表,必须在同一个本地事务里(保证业务与消息原子) */
@Transactional(rollbackFor = Exception.class)
public void createOrderWithOutbox(OrderDTO dto) {
// 1. 业务操作:插入订单
Order order = toEntity(dto);
orderMapper.insert(order);
// 2. 同事务写入本地消息表:如果订单插入失败,这条消息也不会落库
OutboxMessage msg = new OutboxMessage();
msg.setTopic("order.created");
msg.setPayload(JSON.toJSONString(order));
msg.setStatus(0); // 0=待发送
outboxMapper.insert(msg);
log.info("下单成功并写入消息表, orderNo={}", dto.getOrderNo());
}
}
消息表结构:
CREATE TABLE t_outbox_message (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
topic VARCHAR(64) NOT NULL,
payload TEXT NOT NULL,
status TINYINT NOT NULL DEFAULT 0 COMMENT '0=待发送 1=已发送',
retry_count INT NOT NULL DEFAULT 0,
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
KEY idx_status (status)
) ENGINE=InnoDB;
后台投递任务:定时扫描「待发送」消息,投递给
MQ,投递成功才标记为已发送:
@Slf4j
@Component
@RequiredArgsConstructor
public class OutboxRelayTask {
private final OutboxMessageMapper outboxMapper;
private final RabbitTemplate rabbitTemplate;
/** 每 5 秒扫一次待发送消息并投递,保证「业务成功后消息最终一定会被发出去」 */
@Scheduled(fixedDelay = 5000)
public void relay() {
List<OutboxMessage> pending = outboxMapper.selectByStatus(0);
for (OutboxMessage msg : pending) {
try {
rabbitTemplate.convertAndSend("order.exchange", msg.getTopic(), msg.getPayload());
outboxMapper.markSent(msg.getId()); // 投递成功才改状态
} catch (Exception e) {
// 投递失败:增加重试次数,下次扫描继续投。达到上限转人工/告警
log.error("消息投递失败, msgId={}", msg.getId(), e);
outboxMapper.increaseRetry(msg.getId());
}
}
}
}
消费端再配合第 2 章的消息幂等,形成完整的「下单 → 可靠投递 →
幂等消费」闭环。以下是消费端生产级写法(幂等 + 状态机 + 失败重试):
@Slf4j
@Component
@RequiredArgsConstructor
public class OrderCreatedListener {
private final StringRedisTemplate redisTemplate;
private final OrderMapper orderMapper;
@RabbitListener(queues = "order.created.queue")
public void onOrderCreated(OrderCreatedMessage msg) {
String dedupKey = "msg:dedup:order:" + msg.getMessageId();
// 消费端幂等:messageId 占坑,重复投递直接丢弃
Boolean first = redisTemplate.opsForValue()
.setIfAbsent(dedupKey, "1", 24, TimeUnit.HOURS);
if (!Boolean.TRUE.equals(first)) {
log.info("重复消息被丢弃, messageId={}", msg.getMessageId());
return;
}
try {
// 状态机:订单只能从「待支付」流转到「已支付」,重复处理天然幂等
Order order = orderMapper.selectByOrderNo(msg.getOrderNo());
if (order != null && order.getStatus() == OrderStatus.PAID.getCode()) {
log.info("订单已支付,忽略重复消息, orderNo={}", msg.getOrderNo());
return;
}
orderMapper.markPaid(msg.getOrderNo());
log.info("订单状态已更新为已支付, orderNo={}", msg.getOrderNo());
} catch (Exception e) {
// 处理失败删去重 key,让 MQ 重投后能重新处理
redisTemplate.delete(dedupKey);
throw e; // 抛异常 = 不 ACK,MQ 会重新投递
}
}
}
为什么本地消息表能保证最终一致:业务和消息在同一个事务里原子落库,所以「业务成功但消息丢失」的情况不会发生;消息表里的消息由后台任务反复重试投递,只要
MQ
最终可用,消息最终一定能发出去;消费端再配合消息幂等,就能把「下单成功」这个事实可靠地传播到库存、积分等下游,最终达到一致。
一个容易踩的坑:「本地事务提交」和「消息投递」之间没有天然的原子性。如果后台任务在业务事务提交「之前」就把消息发出去了(极端时序),消费端可能查到「订单还不存在」。解决:后台任务只扫描「创建时间早于当前时间几秒」的消息,或消费端做「查不到就重试/延迟消费」的兜底。
5.4 TCC:Try / Confirm / Cancel
TCC
把每个业务操作拆成三个接口,适合需要精确控制的资金类场景:
- Try:预留资源(冻结库存、冻结余额),不做真实扣减;
- Confirm:确认执行(真正扣减);
- Cancel:回滚(释放预留)。
以「下单冻结库存」为例,三个接口的语义:
public interface TccAction {
/** Try:冻结库存,不真正扣减;幂等(同一笔事务重复 Try 返回成功) */
boolean tryFreeze(String tccId, Long productId, int quantity);
/** Confirm:真正扣减库存;幂等 */
boolean confirm(String tccId);
/** Cancel:释放冻结的库存;幂等 + 支持空回滚(Try 没执行成功时 Cancel 也要能处理) */
boolean cancel(String tccId);
}
完整实现(用「冻结表」记录 Try 阶段预留的资源,Confirm/Cancel 都按
tccId 幂等):
@Slf4j
@Service
@RequiredArgsConstructor
public class StockTccAction implements TccAction {
private final StockMapper stockMapper;
private final FreezeMapper freezeMapper;
@Override
@Transactional(rollbackFor = Exception.class)
public boolean tryFreeze(String tccId, Long productId, int quantity) {
// 幂等:同一 tccId 的 Try 已执行过,直接返回成功
if (freezeMapper.existsByTccId(tccId)) {
return true;
}
// 预留资源:把库存从「可用」挪到「冻结」(不真正扣减,只是标记)
int rows = stockMapper.freeze(productId, quantity); // where stock >= quantity,stock -= quantity 冻结到冻结表
if (rows <= 0) {
return false; // 预留失败(库存不足)
}
freezeMapper.insert(tccId, productId, quantity); // 记录冻结明细
return true;
}
@Override
@Transactional(rollbackFor = Exception.class)
public boolean confirm(String tccId) {
// 幂等:已确认过就跳过
Freeze freeze = freezeMapper.selectByTccId(tccId);
if (freeze == null || freeze.getStatus() == FreezeStatus.CONFIRMED.getCode()) {
return true;
}
// 确认:真正扣减(把冻结的库存正式减掉)
stockMapper.commitFreeze(freeze.getProductId(), freeze.getQuantity());
freezeMapper.markConfirmed(tccId);
return true;
}
@Override
@Transactional(rollbackFor = Exception.class)
public boolean cancel(String tccId) {
// 幂等 + 空回滚:Try 可能没执行成功,Cancel 时冻结记录不存在,直接返回成功
Freeze freeze = freezeMapper.selectByTccId(tccId);
if (freeze == null || freeze.getStatus() == FreezeStatus.CANCELLED.getCode()) {
return true; // 空回滚:没有可释放的资源
}
// 回滚:释放冻结(把冻结的库存退回可用)
stockMapper.releaseFreeze(freeze.getProductId(), freeze.getQuantity());
freezeMapper.markCancelled(tccId);
return true;
}
}
TCC
的优点是「不依赖本地事务也能精确控制」,缺点极其明显:每个业务接口要写三份逻辑,开发量和维护成本是普通方案的三倍,且
Cancel 的幂等和空回滚(Try 没执行就被
Cancel)都是坑。只有资金、支付这类「一分钱都不能错」的场景才值得上
TCC。
5.5 Saga:长事务的补偿编排
Saga 把一个长事务拆成一串「本地事务 +
补偿事务」:T1, T2, T3...
顺序执行,任一步失败,就逆序执行对应的补偿 C1, C2, C3
把前面的操作撤销。适合跨多个服务的长流程(如订单 → 支付 → 发货 →
签收),通常用编排(中心协调器)或编排(事件驱动)两种方式实现。
// Saga 思想伪代码:本地事务 + 补偿动作,由编排器协调
public void sagaOrder() {
try {
createOrder(); // T1:创建订单
deductStock(); // T2:扣库存
pay(); // T3:支付
} catch (Exception e) {
// 逆序补偿:撤销前面已成功的步骤
cancelPay(); // C3
restoreStock(); // C2
cancelOrder(); // C1
}
}
Saga 和 TCC 的关键区别:TCC 的补偿动作在 Try
阶段就「预留」好了资源,回滚是「释放预留」;Saga
的补偿是「事后撤销已执行的操作」。Saga
更通用(不要求每个接口改写成三阶段),但要处理补偿动作本身失败、补偿动作的幂等等问题。
5.6 2PC /
XA:强一致的两阶段提交
2PC 是最早的分布式事务协议,把提交拆成两个阶段:
- Prepare(准备阶段):协调者向所有参与者发「prepare」,每个参与者执行本地事务但不提交,返回「可以提交」或「不能提交」;
- Commit /
Rollback(提交阶段):所有参与者都「可以提交」时,协调者发「commit」让所有人提交;任一参与者「不能提交」,协调者发「rollback」让所有人回滚。
Spring Boot 里可以用 JTA(如 Atomikos、Narayana)实现
2PC,但生产上很少用,原因有二:
- 同步阻塞:准备阶段所有资源被锁住,直到提交完成,吞吐极低;
- 协调者单点:协调者宕机会让参与者「卡在
prepare」状态,需要额外机制恢复。
所以 2PC
只适合「同一系统内跨多个同类型数据库、强一致必须、并发量不大」的场景(如财务系统跨两个
MySQL 库),不适合高并发电商。
5.7
事务消息:本地消息表的「进阶版」
本地消息表有个小缺陷:业务和消息虽然同库同事务,但消息的「投递」依赖后台定时任务,有秒级延迟。RocketMQ
的事务消息把「业务 + 消息」的原子性从应用层下沉到 MQ
层,投递更及时:
- 生产者发送「半消息」(消费者暂时收不到);
- 执行本地事务(下单);
- 本地事务成功,提交半消息(消费者可收到);失败则回滚半消息(消费者永远收不到);
- 若生产者宕机没来得及提交/回滚,MQ
会回查本地事务状态,决定提交还是回滚。
它的本质和本地消息表一样,都是「业务与消息原子绑定」,只是把「回查状态」的责任交给了
MQ。RabbitMQ 没有原生事务消息,通常用本地消息表或「发送方确认(publisher
confirm)+ 补偿」来实现类似效果。
5.8 选型决策树
把四种方案放到一个决策流程里,遇到「跨服务要保证一致」时按顺序问自己:
跨服务操作需要事务吗?
├── 否 → 用消息 + 幂等即可(不是所有跨服务都需要「事务」,很多只是「通知」)
└── 是 → 能接受几秒不一致吗?
├── 不能(强一致必须)→ 2PC/XA(但要接受低吞吐、单点风险)
└── 能 → 业务有精确的资金语义吗?
├── 有(钱、库存精确扣减)→ TCC
└── 没有 → 流程长、步骤多吗?
├── 是 → Saga
└── 否 → 本地消息表(首选,性价比最高)
这个决策树的核心精神:先判断「到底需不需要事务」,再判断「能不能最终一致」,最后才选具体方案。绝大多数场景会落在「本地消息表」,而它真正依赖的,恰恰是第
2
章的幂等——这也是本篇把「幂等」放在「分布式事务」之前的用意:最终一致性的可靠性,一半在投递,一半在幂等。
本章小结:分布式事务的真相是「分布式下没有免费的一致」。2PC
用性能和可用性换强一致,本地消息表/TCC/Saga 用「补偿 + 重试 +
幂等」换最终一致。选型口诀:能用本地消息表就用本地消息表,涉及钱才上
TCC,长流程用 Saga,2PC 能不用就不用。
第 6 章:DDD 与代码质量
6.1 什么是 DDD,解决什么问题
领域驱动设计(Domain-Driven Design,DDD)
是一套应对复杂业务的建模与组织代码的方法论。它要解决的核心痛点是:业务越复杂,代码越混乱。
回想你见过的「贫血模型」:Controller
塞满业务逻辑,Service 上千行,Entity 只是一堆
getter/setter,业务规则散落在各处。当业务规则足够复杂(下单要校验库存、计算优惠、分摊运费、处理各种活动叠加),这种写法会让「改一处崩三处」成为日常。
对比两种模型的本质区别:
- 贫血模型(传统三层):
Order只有字段和
getter/setter,业务规则(怎么算金额、怎么校验状态)全部写在
OrderService里。数据和行为分离,规则一多就散落各处。 - 充血模型(DDD):
Order
自己携带业务规则(pay()、addItem()),Service
只做编排。数据和行为内聚,规则集中、易测、易改。
DDD
的主张:让代码结构反映业务结构,把业务规则集中到「领域层」,用领域模型承载业务逻辑。它不改变你的技术栈,改变的是「代码怎么组织」。
6.2 四个核心概念
(1)限界上下文(Bounded Context)
一个「上下文」就是一个业务边界,边界内有一套独立、自治的领域模型和术语。比如电商系统可以划分为:商品上下文、订单上下文、库存上下文、用户上下文、支付上下文。同一个词在不同上下文里含义可能完全不同——「订单」在订单上下文里是「待处理的交易」,在财务上下文里是「一笔应收款」。
限界上下文的价值:给混乱的业务划边界。边界内强一致(用本地事务),边界之间用领域事件做最终一致(第
5
章)。它也是微服务拆分的重要依据——理想情况下一个限界上下文对应一个微服务。
(2)实体(Entity)与值对象(Value Object)
- 实体:有唯一标识(ID),即使属性变了,它还是「同一个」对象。如「用户」「订单」——订单号不变它还是同一个订单。
- 值对象:没有唯一标识,由属性定义,且不可变。如「地址」「金额」「收货人」。两个「地址」只要城市/街道/门牌都一样,就是同一个地址,你不需要给它一个
ID。
判断标准一句话:看它有没有「身份」。有
ID、需要跟踪生命周期的是实体;没
ID、用值来定义、用完即弃的是值对象。
值对象的生产级写法(金额 Money:不可变 + 值相等性 + 运算内聚):
/** 值对象:金额。不可变、无标识、靠「值」比较相等,封装金额运算规则 */
public final class Money {
private final BigDecimal amount;
private final String currency;
public Money(BigDecimal amount, String currency) {
if (amount == null || amount.compareTo(BigDecimal.ZERO) < 0) {
throw new BizException(ErrorCode.PARAM_ERROR.getCode(), "金额不能为空或为负");
}
this.amount = amount;
this.currency = currency;
}
public Money add(Money other) {
if (!this.currency.equals(other.currency)) {
throw new BizException(ErrorCode.PARAM_ERROR.getCode(), "币种不一致不能相加");
}
return new Money(this.amount.add(other.amount), this.currency);
}
// 值对象重写 equals/hashCode:按「值」比较,而不是按引用
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (!(o instanceof Money other)) return false;
return amount.compareTo(other.amount) == 0 && currency.equals(other.currency);
}
@Override
public int hashCode() {
return Objects.hash(amount, currency);
}
}
(3)聚合(Aggregate)与聚合根(Aggregate Root)
聚合是一组「必须一起保证一致性」的实体和值对象的集合,对外只通过聚合根访问。经典例子:订单聚合,聚合根是
Order,它包含多个 OrderItem(订单明细)和
Address(收货地址值对象)。外部不能直接改
OrderItem,必须通过 Order 来改。
聚合的价值:把一致性边界圈出来。聚合内的修改要原子完成(一个事务),聚合之间通过
ID 引用 +
领域事件协作,不直接持有对方的引用。这个「聚合内强一致、聚合间最终一致」的划分,直接对应第
5 章的一致性设计。
(4)领域事件(Domain Event)
领域里「已经发生的重要事实」,用过去式命名,如
OrderPaid(订单已支付)、StockDeducted(库存已扣减)。它把「一个聚合内部的变化」广播给「其他关心这件事的上下文」,是限界上下文之间解耦的桥梁,也天然对应第
5 章的最终一致性和第 4 章的异步削峰。
领域事件的实现(聚合根内发布,Spring 事件驱动):
/** 领域事件:订单已支付(用过去式命名,表达「已经发生的事实」) */
public record OrderPaidEvent(Long orderId, String orderNo, BigDecimal amount) {
}
// 在聚合根内发布领域事件
public void pay() {
if (this.status != OrderStatus.PENDING) {
throw new BizException(ErrorCode.PARAM_ERROR.getCode(), "订单状态不允许支付");
}
this.status = OrderStatus.PAID;
// 发布领域事件:其他上下文(积分、通知)订阅后各自处理
registerEvent(new OrderPaidEvent(this.id, this.orderNo, this.totalAmount));
}
6.3 DDD 分层架构
DDD 的分层是「四层」,和传统的 Controller/Service/Mapper
三层最关键的区别是:业务规则全部沉到 domain
层,成为系统的核心,其他层都围绕它转。
interfaces(接口层) → Controller、DTO/VO 转换、参数校验,只做「翻译」,不含业务
↓ 调用
application(应用层) → 用例编排:协调领域对象完成一个用例,薄薄一层,不含业务规则
↓ 调用
domain(领域层)★核心 → Entity、Value Object、Aggregate、Domain Service、Domain Event,业务规则全在这
↓ 依赖接口(面向接口,不依赖实现)
infrastructure(基础设施层)→ Repository 实现、Mapper、Redis、MQ,实现 domain 定义的接口
依赖方向是单向的:interfaces → application →
domain,domain 不依赖任何外层(依赖倒置:domain 定义
OrderRepository 接口,infrastructure 提供实现)。
用一个下单用例演示 DDD
分层。领域层——订单聚合根(业务规则集中于此):
/**
* 订单聚合根:封装订单的业务规则。
* 所有对订单明细、状态的修改都必须通过它,外部不能直接改明细。
*/
public class Order {
private Long id;
private String orderNo;
private Long userId;
private List<OrderItem> items = new ArrayList<>();
private BigDecimal totalAmount = BigDecimal.ZERO;
private OrderStatus status;
/** 工厂方法:创建订单,封装「订单号生成 + 状态初始化 + 金额计算」规则 */
public static Order create(Long userId, List<OrderItem> items) {
Order order = new Order();
order.userId = userId;
order.items = new ArrayList<>(items);
order.status = OrderStatus.PENDING;
order.orderNo = generateOrderNo();
order.recalculateAmount(); // 创建时算一次总金额
return order;
}
/** 领域方法:支付。把「已支付」的状态流转规则放在领域对象内部,而不是散落在 Service */
public void pay() {
if (this.status != OrderStatus.PENDING) {
throw new BizException(ErrorCode.PARAM_ERROR.getCode(), "订单状态不允许支付");
}
this.status = OrderStatus.PAID;
// 这里可以发布领域事件 OrderPaid,供其他上下文订阅
}
/** 领域方法:添加明细并重算金额,保证「金额永远由明细推导,不会算错」 */
public void addItem(OrderItem item) {
if (this.status != OrderStatus.PENDING) {
throw new BizException(ErrorCode.PARAM_ERROR.getCode(), "订单已提交,不能修改明细");
}
this.items.add(item);
recalculateAmount();
}
/** 金额计算是核心业务规则,集中在领域层,任何入口下单都复用同一套算法 */
private void recalculateAmount() {
this.totalAmount = items.stream()
.map(i -> i.getPrice().multiply(BigDecimal.valueOf(i.getQuantity())))
.reduce(BigDecimal.ZERO, BigDecimal::add);
}
// getter 省略;值对象/构造细节见完整项目
}
领域层——Repository 接口(只定义接口,不依赖
MyBatis-Plus 具体实现):
/** 领域层定义仓储接口:面向领域对象,屏蔽持久化细节 */
public interface OrderRepository {
Order findById(Long id);
void save(Order order);
}
应用层——用例编排(薄薄一层,只协调,不含业务规则):
@Slf4j
@Service
@RequiredArgsConstructor
public class OrderApplicationService {
private final OrderRepository orderRepository;
/** 应用服务:编排「创建订单」用例,业务规则在 Order 聚合根内部 */
public void createOrder(CreateOrderCommand command) {
// 1. 组装领域对象(业务规则:金额计算、状态初始化在 Order.create 里)
Order order = Order.create(command.getUserId(), command.getItems());
// 2. 持久化(通过仓储接口,具体实现是 MyBatis-Plus)
orderRepository.save(order);
log.info("订单创建成功, orderNo={}", order.getOrderNo());
}
}
基础设施层——Repository 实现:
@Repository
@RequiredArgsConstructor
public class OrderRepositoryImpl implements OrderRepository {
private final OrderMapper orderMapper;
@Override
public Order findById(Long id) {
return orderMapper.selectById(id);
}
@Override
public void save(Order order) {
orderMapper.insert(order);
}
}
接口层——Controller 只做翻译:
@RestController
@RequiredArgsConstructor
public class OrderController {
private final OrderApplicationService orderApplicationService;
@PostMapping("/orders")
public ApiResult<Void> create(@Valid @RequestBody CreateOrderCommand command) {
orderApplicationService.createOrder(command);
return ApiResult.ok(null);
}
}
一个常见的疑问:应用服务(Application
Service)和领域服务(Domain
Service)有什么区别?应用服务是「用例编排器」,它不持有业务规则,只负责「调谁、按什么顺序调、要不要开事务」;领域服务是「跨聚合的业务规则」,当一段业务规则不属于任何一个聚合根、又必须协调多个聚合时,才放进领域服务(如「转账」同时涉及两个账户聚合,规则不属于任何一个账户,就放
TransferDomainService)。判断标准:能在聚合根里写的规则,就不要建领域服务;应用服务永远不应该写业务规则。
用一个「转账」领域服务把这条判断讲透——「转账」同时改两个账户,规则既不属于「转出账户」也不属于「转入账户」,所以放领域服务:
/** 领域服务:跨聚合的转账规则。规则不属于任何一个账户聚合根,故独立成领域服务 */
public class TransferDomainService {
private final AccountRepository accountRepository;
/** 转账:校验余额 → 转出 → 转入,规则协调两个账户聚合 */
public void transfer(Long fromAccountId, Long toAccountId, Money amount) {
Account from = accountRepository.findById(fromAccountId);
Account to = accountRepository.findById(toAccountId);
// 业务规则:余额校验(规则不放在任何单一账户内部,因为它是「跨账户」的约束)
if (from.getBalance().compareTo(amount) < 0) {
throw new BizException(ErrorCode.PARAM_ERROR.getCode(), "余额不足");
}
from.debit(amount); // 账户聚合根自己的领域方法:扣款
to.credit(amount); // 账户聚合根自己的领域方法:入账
accountRepository.save(from);
accountRepository.save(to);
}
}
而「扣款/入账」这种属于单个账户的规则,则放在 Account
聚合根内部:
/** 账户聚合根:扣款规则属于单个账户,写在聚合根内部 */
public class Account {
private Long id;
private Money balance;
public void debit(Money amount) {
if (balance.compareTo(amount) < 0) {
throw new BizException(ErrorCode.PARAM_ERROR.getCode(), "余额不足");
}
this.balance = this.balance.subtract(amount);
}
public void credit(Money amount) {
this.balance = this.balance.add(amount);
}
}
6.4.1
重构:有测试保护,小步快跑
DDD
分层重构不是「推倒重写」,而是「小步快跑」。一个可落地的重构路径:
- 先补测试:给现有
OrderService
的核心用例写测试(阶段 8),建立安全网; - 提取领域规则:把
OrderService
里「算金额、校验状态流转」的逻辑,逐步搬进Order
聚合根的方法里,每搬一步跑一次测试; - 引入仓储接口:把
OrderMapper
的直接调用,替换成OrderRepository
接口(依赖倒置),实现放基础设施层; - 拆应用服务:让
OrderService
退化成「只编排、不写规则」的薄应用服务。
每一步都是「小改动 + 测试绿」,而不是一次性大爆炸式重写——这是 DDD
落地最务实、也最不容易翻车的姿势。
6.4 SOLID
与设计模式(代码质量的另一面)
DDD 解决「业务建模」,SOLID
解决「代码设计」,两者是代码质量的两根支柱。
SOLID 五原则:
| 原则 | 一句话 | 反面典型 |
|---|---|---|
| 单一职责(SRP) | 一个类只为一个理由改变 | 一个 Service 又管订单又管库存又管短信 |
| 开闭原则(OCP) | 对扩展开放,对修改关闭 | 加一种支付方式要改 if-else 大串联 |
| 里氏替换(LSP) | 子类能替换父类而不出错 | 子类重写方法后违反父类约定 |
| 接口隔离(ISP) | 接口小而专,不要大而全 | 一个接口 20 个方法,实现类全都要实现 |
| 依赖倒置(DIP) | 依赖抽象,不依赖具体实现 | Service 直接 new 具体类 |
设计模式里,最值得在后端掌握的是策略模式(消除
if-else) 和 模板方法(统一流程骨架)。
策略模式消除 if-else——支付方式为例。坏味道写法:
// 反例:每加一种支付方式,就要加一个 else-if,OCP 被破坏
public void pay(String payType, BigDecimal amount) {
if ("alipay".equals(payType)) {
alipayPay(amount);
} else if ("wechat".equals(payType)) {
wechatPay(amount);
} else if ("unionpay".equals(payType)) {
unionpayPay(amount);
}
}
策略模式重构:
/** 支付策略接口 */
public interface PayStrategy {
/** 返回该策略支持的支付类型,用于路由 */
String type();
void pay(BigDecimal amount);
}
@Component
public class AlipayStrategy implements PayStrategy {
@Override public String type() { return "alipay"; }
@Override public void pay(BigDecimal amount) { log.info("支付宝支付 {}", amount); }
}
@Component
public class WechatStrategy implements PayStrategy {
@Override public String type() { return "wechat"; }
@Override public void pay(BigDecimal amount) { log.info("微信支付 {}", amount); }
}
/** 策略工厂:用 Map 按 type 路由,新增支付方式只需新增一个类,不改任何老代码 */
@Component
@RequiredArgsConstructor
public class PayStrategyFactory {
private final Map<String, PayStrategy> strategyMap;
public PayStrategyFactory(List<PayStrategy> strategies) {
// 启动时把所有策略按 type 收集进 Map
this.strategyMap = strategies.stream()
.collect(Collectors.toMap(PayStrategy::type, s -> s));
}
public PayStrategy get(String type) {
PayStrategy strategy = strategyMap.get(type);
if (strategy == null) {
throw new BizException(ErrorCode.PARAM_ERROR.getCode(), "不支持的支付方式: " + type);
}
return strategy;
}
}
模板方法模式——把「固定流程」抽到父类,让子类只实现变化的部分。典型场景:所有支付都要「校验
→ 扣款 → 落账 → 发通知」,但每种支付方式的「扣款」实现不同:
/** 模板方法:定义支付流程骨架,扣款细节交给子类 */
public abstract class AbstractPayTemplate {
/** 模板方法:固定流程,子类不得修改(final 防止被覆写) */
public final void pay(BigDecimal amount) {
validate(amount); // 1. 校验(公共)
doPay(amount); // 2. 扣款(子类实现,变化点)
record(amount); // 3. 落账(公共)
notifyUser(amount); // 4. 通知(公共)
}
protected void validate(BigDecimal amount) {
if (amount.compareTo(BigDecimal.ZERO) <= 0) {
throw new BizException(ErrorCode.PARAM_ERROR.getCode(), "金额必须大于 0");
}
}
/** 抽象方法:子类实现具体的扣款逻辑 */
protected abstract void doPay(BigDecimal amount);
private void record(BigDecimal amount) { log.info("记录支付流水 {}", amount); }
private void notifyUser(BigDecimal amount) { log.info("通知用户支付成功 {}", amount); }
}
6.5 DDD
适用场景:什么时候该用,什么时候别用
这是本章(也是面试)最务实的一问。DDD
不是「更高级」,而是「更适合复杂业务」,用错场景是负担。
该用 DDD(业务复杂、规则密集、会长期演进):
- 电商核心域(订单、支付、库存、促销),规则多且互相叠加;
- 金融/保险/风控,业务规则复杂且高价值;
- 有多个限界上下文、需要清晰边界的大型单体或微服务。
不该用 DDD(简单 CRUD,硬上就是杀鸡用牛刀):
- 管理后台、配置系统、字典维护——就是增删改查,没有复杂规则;
- 报表、数据查询类系统——没有业务规则,只有数据展示;
- 团队小、迭代快、业务还在探索期——DDD 的建模成本会拖慢试错速度。
判断口诀:看「业务规则密度」。规则越多、越复杂、越容易改,DDD
收益越大;反过来,如果大部分代码就是 selectById /
insert,三层架构已经足够,硬套 DDD
只会让团队「为抽象而抽象」,复杂度爆炸。
本章小结:DDD
的核心是「用领域模型承载业务规则」,四个概念(限界上下文、实体/值对象、聚合、领域事件)分别回答「边界在哪、对象有没有身份、谁保证一致、上下文怎么通信」四个问题。它必须和
SOLID、设计模式配合,并且只在业务足够复杂时才值得投入。
5.
生产级实战项目:抢购系统(分布式锁 + 幂等 + 乐观锁 + 限流)
5.1 项目目标
实现一个电商「抢购」功能,串起本篇至少 80%
的知识点:分布式锁(防超卖)+ 接口幂等(防重复下单)+
乐观锁扣库存(防并发)+ 限流(防刷子)+ 分布式
ID(雪花主键)。
- 场景:10 万用户同时抢 1000 件商品。
- 链路:
限流→幂等校验→
分布式锁→乐观锁扣库存→
创建订单。
5.2 目录结构
com.example.demo
├── controller
│ └── SeckillController # 接口层:限流 + 调用服务
├── service
│ └── SeckillService # 业务层:编排锁 + 幂等 + 扣库存
├── mapper
│ ├── StockMapper # 库存数据访问(乐观锁 SQL)
│ └── OrderMapper # 订单数据访问
├── entity
│ ├── Stock # 库存实体
│ └── Order # 订单实体(雪花 ID)
├── dto
│ └── SeckillDTO # 入参:requestNo + userId + productId
├── common
│ ├── ApiResult # 统一响应(复用)
│ ├── ErrorCode # 统一错误码(复用 + 扩展)
│ └── RedisRateLimiter # 限流器
├── exception
│ ├── BizException # 业务异常(复用)
│ └── GlobalExceptionHandler # 全局异常(复用)
└── config
└── RedissonConfig # Redisson 配置
5.3 数据库表结构
-- 库存表:stock 用乐观锁保护,num > 0 才允许扣减
CREATE TABLE t_stock (
product_id BIGINT PRIMARY KEY,
stock INT NOT NULL,
version INT NOT NULL DEFAULT 0, -- 乐观锁版本号
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
) ENGINE=InnoDB;
-- 订单表:order_no 唯一索引兜底幂等,主键用雪花算法
CREATE TABLE t_order (
id BIGINT PRIMARY KEY, -- 雪花算法生成
order_no VARCHAR(64) NOT NULL,
user_id BIGINT NOT NULL,
product_id BIGINT NOT NULL,
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
UNIQUE KEY uk_order_no (order_no),
UNIQUE KEY uk_user_product (user_id, product_id) -- 一人一单,幂等兜底
) ENGINE=InnoDB;
初始化数据:
INSERT INTO t_stock (product_id, stock, version) VALUES (1001, 1000, 0);
5.4 核心代码
库存扣减(乐观锁)——StockMapper:
@Mapper
public interface StockMapper {
/**
* 乐观锁扣减库存:where stock >= quantity 保证「库存够才扣」,天然防超卖。
* 返回受影响行数:0 表示库存不足或被并发抢走。
*/
@Update("UPDATE t_stock SET stock = stock - #{quantity}, version = version + 1 " +
"WHERE product_id = #{productId} AND stock >= #{quantity}")
int deduct(@Param("productId") Long productId, @Param("quantity") int quantity);
}
订单实体(雪花 ID + 幂等)——Order:
@Data
@TableName("t_order")
public class Order {
/** ASSIGN_ID:MyBatis-Plus 内置雪花算法,分库分表后全局唯一 */
@TableId(type = IdType.ASSIGN_ID)
private Long id;
private String orderNo;
private Long userId;
private Long productId;
private LocalDateTime createTime;
}
入参 DTO(带
requestNo,幂等)——SeckillDTO:
@Data
public class SeckillDTO {
/** 请求唯一标识:前端生成,同一笔抢购重试必须传同一个值 */
@NotBlank(message = "requestNo 不能为空")
private String requestNo;
@NotNull(message = "用户 id 不能为空")
private Long userId;
@NotNull(message = "商品 id 不能为空")
private Long productId;
}
抢购核心服务(锁 + 幂等 +
乐观锁三层组合)——SeckillService:
@Slf4j
@Service
@RequiredArgsConstructor
public class SeckillService {
private final RedissonClient redissonClient;
private final StringRedisTemplate redisTemplate;
private final StockMapper stockMapper;
private final OrderMapper orderMapper;
private static final String IDEMPOTENT_PREFIX = "idempotent:seckill:";
private static final String LOCK_PREFIX = "lock:seckill:";
/**
* 抢购主流程。三层防护,缺一不可:
* 1. 幂等(Redis setNX + requestNo):防「同一个请求重试」重复下单;
* 2. 分布式锁(Redisson):防「多实例并发」同时操作同一商品;
* 3. 乐观锁(where stock >= 1):防「锁内」的最终兜底,杜绝超卖。
*/
public void seckill(SeckillDTO dto) {
// ---- 第一层:接口幂等(requestNo 占坑,防重复提交)----
String idempotentKey = IDEMPOTENT_PREFIX + dto.getRequestNo();
Boolean first = redisTemplate.opsForValue()
.setIfAbsent(idempotentKey, "1", 10, TimeUnit.MINUTES);
if (!Boolean.TRUE.equals(first)) {
log.info("重复抢购请求被拦截, requestNo={}", dto.getRequestNo());
throw new BizException(ErrorCode.DUPLICATE_SUBMIT);
}
String lockKey = LOCK_PREFIX + dto.getProductId();
RLock lock = redissonClient.getLock(lockKey);
boolean locked = false;
try {
// ---- 第二层:分布式锁(多实例互斥,防止同时扣同一商品)----
locked = lock.tryLock(2, TimeUnit.SECONDS);
if (!locked) {
log.warn("获取分布式锁失败, productId={}", dto.getProductId());
throw new BizException(ErrorCode.TOO_MANY_REQUESTS);
}
// ---- 第三层:乐观锁扣库存(最终兜底,库存不足返回 0)----
int rows = stockMapper.deduct(dto.getProductId(), 1);
if (rows <= 0) {
log.info("库存不足, productId={}", dto.getProductId());
throw new BizException(ErrorCode.SOLD_OUT);
}
// ---- 创建订单(order_no 唯一索引 + user_product 唯一索引双重幂等)----
Order order = new Order();
order.setOrderNo("SK" + dto.getRequestNo());
order.setUserId(dto.getUserId());
order.setProductId(dto.getProductId());
order.setCreateTime(LocalDateTime.now());
try {
orderMapper.insert(order);
} catch (DuplicateKeyException e) {
// 唯一索引冲突:该用户已抢到,视为重复下单,幂等吞掉
log.info("用户重复下单被唯一索引拦截, userId={}, productId={}",
dto.getUserId(), dto.getProductId());
throw new BizException(ErrorCode.DUPLICATE_SUBMIT);
}
log.info("抢购成功, userId={}, productId={}, orderNo={}",
dto.getUserId(), dto.getProductId(), order.getOrderNo());
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new BizException(ErrorCode.SYSTEM_ERROR);
} finally {
// 释放锁:只释放当前线程持有的锁,防止误删
if (locked && lock.isHeldByCurrentThread()) {
lock.unlock();
}
// 注意:业务失败时这里不删幂等 key,因为抢购失败也「已经处理过这次请求」,
// 删掉反而会让用户拿着同一个 requestNo 反复刷。真正需要「失败可重试」的场景才删。
}
}
}
接口层(限流 +
调用)——SeckillController:
@RestController
@RequiredArgsConstructor
public class SeckillController {
private final RedisRateLimiter rateLimiter;
private final SeckillService seckillService;
@PostMapping("/seckill")
public ApiResult<Void> seckill(@Valid @RequestBody SeckillDTO dto) {
// 入口限流:同一用户每秒最多 5 次,把刷子挡在业务逻辑之外
if (!rateLimiter.tryAcquire("rate:seckill:" + dto.getUserId(), 5, 1)) {
throw new BizException(ErrorCode.TOO_MANY_REQUESTS);
}
seckillService.seckill(dto);
return ApiResult.ok(null);
}
}
5.5 运行步骤
# 1. 启动依赖
docker run -d --name redis -p 6379:6379 redis:7
docker run -d --name mysql -p 3306:3306 -e MYSQL_ROOT_PASSWORD=root mysql:8.0
# 2. 初始化表结构与库存(执行 5.3 的 SQL)
# 3. 启动应用
mvn spring-boot:run
# 4. 模拟并发抢购(Apache Bench 压测,1000 并发请求)
# 注意:每个请求的 requestNo 要不同(模拟不同用户的不同请求)
ab -n 2000 -c 200 -p seckill.json -T application/json http://localhost:8080/seckill
seckill.json(压测请求体示例,requestNo
用脚本生成不同值):
{"requestNo":"req-0001","userId":1,"productId":1001}
验证结果:压测结束后,t_order
表的行数应该等于 t_stock.stock
被扣减的数量,且不会超过初始库存
1000,也不会出现重复订单——这就是「锁 + 幂等 +
乐观锁」三层防护共同作用的结果。
5.6 为什么是三层,而不是一层
这是本项目最重要的一个理解点。单独用任何一层,在极端场景下都有漏洞:
- 只用幂等:挡得住「同一个 requestNo
重复」,挡不住「不同用户并发抢同一商品」导致的超卖; - 只用分布式锁:挡得住并发扣库存,但锁一旦因为 GC
停顿等被释放,锁内的乐观锁仍能兜底;而「用户点了两次生成两个
requestNo」这种重复,锁是挡不住的; - 只用乐观锁:挡得住超卖(
where stock >= 1),但库存不足时大量请求直接打到数据库,缺少锁的「串行化」保护,也挡不住重复下单。
三层各管一段:幂等管「同一请求重复」,锁管「同一时刻互斥」,乐观锁管「最终不超卖」。生产级的高并发系统,从来不是靠一个「完美的锁」解决的,而是靠分层防御。
6. 常见坑与排错指南
| 坑/现象 | 原因 | 解决方案 |
|---|---|---|
| 分布式锁不加过期时间,进程崩溃后所有请求永久阻塞 | SET NX 后没有 EX,锁不会自动释放 |
加锁必须用 SET key value NX EX seconds 原子命令,或Redisson 默认租期 |
| 锁被「别人」释放,出现两个线程同时进临界区 | 解锁只 DEL key,没校验这把锁是不是自己加的 |
解锁用 Lua 脚本校验 value == requestId 后再DEL |
| 业务执行超过锁过期时间,锁提前释放导致并发 | 锁超时时间定得太短,又没有续期 | 用 Redisson 看门狗自动续期;或显式 leaseTime 必须 > 业务最长执行时间 |
| 同一请求重试导致重复扣款/重复下单 | 接口非幂等,网络重试/用户双击被重复执行 | 请求号 + Redis setNX 占坑 + 唯一索引/状态机兜底 |
| 雪花算法生成重复 ID | 系统时钟回拨(NTP 校时/手动改时间/虚拟机迁移) | 回拨小则自旋等待追平,回拨超阈值则拒绝并告警;或切换备用 workerId |
| 主键用 UUID 后插入越来越慢 | UUID 无序,InnoDB 聚簇索引频繁页分裂 | 主键改用雪花算法(趋势递增) |
| 分库分表后按非分片键查询全表扫描 | 分片键选错,查询不命中分片 | 选「覆盖绝大多数查询」的字段作分片键;非分片键查询引入 ES 或基因法 |
| 简单 CRUD 系统硬上 DDD,代码量翻倍还更慢 | DDD 建模成本高,业务规则密度低时得不偿失 | 按业务复杂度选型,简单 CRUD 用三层架构即可 |
| 分库分表不规划数据迁移,切流后数据错乱 | 存量数据没按分片规则迁移就上线 | 提前设计「双写 → 迁移 → 灰度切流 → 校验」迁移方案 |
| 本地消息表消费端不幂等,消息重投导致重复处理 | 消费端处理完才 ACK,重投后重复执行 | 消费端用 messageId 去重 + 状态机/唯一索引幂等 |
8. 总结与延伸阅读
8.1 本篇总结
本篇是「从工程师到架构师」的收官,六个章节其实在回答同一个问题:当系统从「单机、单库、单表」走向「多实例、多库、多表」后,如何继续保持正确、可用、可维护。分布式锁解决「多实例下的互斥」(重点是防死锁、防误删),幂等解决「重试下的正确性」(重点是唯一索引、状态机、请求号三层叠加),分布式
ID
解决「拆分后的主键唯一」(重点是雪花算法与时钟回拨),高并发架构解决「流量下的吞吐」(重点是缓存/异步/限流/分库分表的组合拳),分布式事务解决「跨库下的一致性」(重点是
CAP/BASE 与最终一致性的落地),DDD
解决「复杂业务下的可维护性」(重点是限界上下文、聚合与「该不该用」的判断)。
贯穿全篇的一条暗线是:所有架构决策都是权衡,没有免费的午餐。锁要防死锁就得加超时,加了超时就得防误删;要最终一致就得接受短暂不一致,要最终一致就得写好补偿和幂等。理解「代价在哪」,比记住「怎么配置」重要得多。
8.2 毕业项目:串联 0–11
全部阶段
学完 0–11
阶段,你已经具备独立搭建一个生产级电商后端的能力。下面这个「毕业项目」是检验你「能否写生产级代码」的最终标准——它把
12 篇文章的知识点全部串联起来:
搭建一个完整的电商后端系统,需要包含:
- 阶段 0–1:环境搭建、Spring Boot 入门、Maven
工程结构; - 阶段
2–3:分层架构(Controller/Service/Mapper/Entity/DTO/VO)、统一响应
ApiResult、全局异常
GlobalExceptionHandler、参数校验(JSR-303)、配置管理; - 阶段 4:MyBatis-Plus
数据访问、事务(传播行为)、Redis 缓存(Cache Aside +
缓存穿透/击穿/雪崩防护); - 阶段 5:用户认证(JWT)+ RBAC 权限控制;
- 阶段 6:RabbitMQ
消息队列、下单异步解耦、死信队列超时关单、消息可靠投递与幂等; - 阶段 7:微服务拆分、Nacos 注册配置、Feign
调用、Gateway 网关、Sentinel 限流熔断; - 阶段 8:单元测试(JUnit 5 + Mockito)+
集成测试(Testcontainers); - 阶段 9:Docker 容器化、K8s
部署、监控告警(Prometheus + Grafana); - 阶段 10:CI/CD 流水线、日志与链路追踪(ELK /
SkyWalking); - 阶段 11(本篇):抢购场景的分布式锁 + 幂等 + 乐观锁
+ 限流 + 分库分表设计 + DDD 分层重构核心域。
毕业项目的验收标准:能扛住模拟并发压测、能追踪一条请求从网关到数据库的完整链路、能在任一服务宕机后自动恢复、代码分层清晰且核心域可维护。做到这些,你就真正从「会写
Spring Boot」走到了「能设计 Spring Boot 系统」。
8.3 延伸阅读
- Redisson
官方文档:https://github.com/redisson/redisson/wiki
—— 分布式锁、看门狗、集群模式的一手资料。 - Apache ShardingSphere
官方文档:https://shardingsphere.apache.org/——
分库分表、读写分离、数据治理的权威参考。 - 《数据密集型应用系统设计》(DDIA):Martin
Kleppmann 著 ——
分布式系统(一致性、复制、分区)的奠基之作,架构师必读。 - 《实现领域驱动设计》:Vaughn Vernon 著 —— DDD
从概念到落地(聚合、限界上下文、事件)的经典。 - 《凤凰架构》:周志明著 ——
从单体到微服务、从强一致到最终一致的演进,中文架构读物里最系统的一本。
文档生成日期:2026-08-14