【阶段 11】架构进阶:分布式系统设计与 DDD

30次阅读
没有评论

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

学完本篇,你能:

  1. 独立实现一个生产级分布式锁,并解释清楚「为什么加锁必须带超时」「为什么解锁必须用
    Lua 校验 value」,能说出 Redisson 看门狗续期的原理,能对比
    Redis、ZooKeeper、数据库三种锁方案。
  2. 为一个下单/扣款接口设计完整的幂等方案,覆盖唯一索引、状态机、请求号、Redis
    setNX、乐观锁五种手段,并能说明「消息幂等」与「接口幂等」的差异。
  3. 手写一个雪花算法 ID 生成器,讲清 64
    位每一位的含义,并能处理时钟回拨问题;能对比
    UUID/雪花/号段/Redis INCR 四类方案的优劣。
  4. 为高并发系统设计缓存 + 异步 + 限流 + 读写分离 +
    分库分表
    的组合方案,说清「什么时候该分库分表、按什么键怎么分」,并能说出分库分表后的三个连锁问题。
  5. 讲清 CAP/BASE
    的区别,并能用本地消息表落地一次最终一致性的分布式事务;能对比本地消息表、TCC、Saga
    三类方案的适用场景。
  6. 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
的公共基础设施(ApiResultErrorCodeBizExceptionGlobalExceptionHandler),不再重新定义。其中
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 为什么需要分布式锁

先回顾一个你早已熟悉的场景。单机下,要防止两个线程同时扣库存,你用
synchronizedReentrantLock 就够了:

// 单机互斥: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)之间,对同一个共享资源(一段缓存、一张表的某一行、一个文件)实现互斥访问。

一个合格的分布式锁,必须同时满足四个条件:

  1. 互斥性:同一时刻只能有一个客户端持有锁(这是锁的底线)。
  2. 不会死锁:即使持有锁的进程崩溃、没来得及释放,锁也必须能自动过期,不能永久卡死。
  3. 不会误删:锁只能由「持有它的那个客户端」释放,不能出现「A
    的锁被 B 释放」。
  4. 高可用:锁服务(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
    秒后锁自动释放,不会永久卡死。

NXEX
必须放在同一条命令里,这是有讲究的。如果拆成两条命令(先
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、下游抖动都会导致):

  1. 线程 A 拿到锁,开始执行,业务跑了 40 秒;
  2. 到第 30 秒,A 的锁过期自动释放了;
  3. 线程 B 在第 31 秒成功拿到锁,开始执行;
  4. 第 40 秒,A 的业务执行完了,调用 DEL lockKey
    想「释放自己的锁」;
  5. 结果 A 把 B 正在持有的锁删掉了;
  6. 线程 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);      // 删除:但这两步之间,锁可能已经过期并被别人拿走
    }
}

GETDEL 之间,锁可能过期、可能被 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 的几个关键点,必须理解清楚:

  1. 看门狗只在 leaseTime == -1
    时生效
    tryLock(3, TimeUnit.SECONDS) 只传了
    waitTime,leaseTime 走默认值 -1,此时 Redisson 会给锁一个默认 30
    秒的租期,并启动看门狗,每 10 秒(租期的 1/3)自动续期一次,直到主动
    unlock。
  2. 如果显式指定了 leaseTime(如
    tryLock(3, 30, TimeUnit.SECONDS)),则不启用看门狗,锁
    30 秒后必过期——适合你明确知道业务一定会在 30 秒内完成的场景。
  3. isHeldByCurrentThread() 的判断很关键:Redisson
    的锁是可重入的,如果 A 线程加了两次锁,第一次的 unlock
    只是把重入计数减一,并没有真正释放。用这个判断可以避免「还没真正释放就提前走
    finally」的误操作。
  4. 可重入: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)),在软件里的含义是:同一个操作执行一次和执行多次,产生的业务结果完全相同

为什么在分布式系统里幂等是「必选项」而不是「加分项」?因为分布式环境下的「重试」无处不在:

  1. 网络超时重试:下单请求发出去,服务端其实已经扣款成功了,但响应在网络里丢了,客户端超时后自动重试,导致重复扣款。
  2. 用户重复点击:手抖连点两下「提交」,或前端按钮没做防抖,同一笔订单提交了两次。
  3. 消息重复投递:RabbitMQ 的
    at-least-once
    投递语义下,消费者宕机后消息会重新投递,你无法保证每条消息只消费一次。
  4. 失败重试框架:你主动引入的重试机制(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());
    }
}

这段代码有两个必须理解的设计决策(也是面试高频考点):

  1. 为什么失败要删
    key
    :如果不删,业务因为参数错误/库存不足失败后,用户改好参数再次提交,会被误判成「重复提交」永远提交不了。删掉
    key 表示「这次尝试作废,允许重试」。
  2. 为什么占坑成功但业务执行到一半进程崩溃:此时 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 要满足:

  1. 全局唯一:任何两张表、任何两个实例生成的 ID
    都不能重复。
  2. 趋势递增:对 InnoDB
    这类聚簇索引,主键有序能显著提升插入性能(避免页分裂),也方便排序和范围查询。
  3. 高性能:高并发下生成 ID
    不能成为瓶颈(要能到每秒几十万甚至百万级)。
  4. 不含敏感信息:自增 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

生产级处理时钟回拨有几种策略:

  1. 短时间回拨等待:回拨在几毫秒内,自旋等待时钟追上来(最小实现已含此逻辑,但只针对「序列号溢出」场景,要扩展成「回拨也等待」)。
  2. 超过阈值拒绝:回拨超过一定毫秒数(如
    5ms),说明不是抖动而是真回拨,此时抛异常或降级。
  3. 备用 workerId:检测到回拨后,临时切换到一个备用机器
    ID 生成,等时钟恢复再切回。
  4. 用外部单调时钟兜底:如从 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 高并发系统的核心手段全景

「高并发」不是某一个技术,而是一组手段的组合。面对流量,架构师手里的牌大致五张,按「离用户由近到远」排列:

  1. 缓存:读多写少的场景,用 Redis 挡住 90%
    的读流量,减轻数据库压力(Cache Aside 模式)。
  2. 异步:写操作走消息队列异步化,把「同步等待」变成「削峰填谷」,让瞬时流量被队列缓冲。
  3. 限流/降级/熔断:在入口处拒绝超额请求(限流),核心链路挂了返回兜底(降级),下游故障快速失败(熔断),保护系统不被冲垮。
  4. 读写分离:主库写、从库读,把读流量分散到多个从库。
  5. 分库分表:当单库单表成为物理瓶颈(数据量、连接数、写入量)时,水平拆分。

前四张牌你其实在前面的阶段已经见过:缓存是阶段 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] = 限流 keyARGV[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

分库分表后有三个「连锁反应」必须提前规划:

  1. 全局 ID:主键不能用自增了,改用第 3
    章的雪花算法。
  2. 数据迁移:存量数据要按分片规则搬到新库新表(用双写
    + 迁移工具 + 灰度切流),这是分库分表最容易翻车的环节。
  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
数据迁移:分库分表最危险的一步

分库分表的技术本身不难,难的是在不影响线上业务的前提下,把存量数据搬到新库新表。标准迁移流程分四步:

  1. 双写:改造写路径,让新数据同时写旧表和新分片(通过
    Canal 等 binlog 同步工具,或应用层双写)。此时以旧表为准。
  2. 历史迁移:用离线任务把存量数据按分片规则批量搬到新库新表,跑完后做数据校验(行数、抽样对账)。
  3. 灰度切流:先切 1%
    流量读新库新表,观察无误后逐步放量到 100%。此时以新分片为准。
  4. 下线旧表:稳定运行一段时间(如一周)后,停止双写,下线旧表。

整个过程中「什么时候以哪边为准」是核心——切流前以旧表为准(新表只是影子),切流后以新分片为准(旧表只读)。这也是为什么「分库分表要提前规划」——如果业务已经跑起来了才想起来分,迁移成本会成倍上升。

本章小结:高并发架构没有银弹,只有「根据瓶颈在哪、选对应的牌」。缓存治读多,异步治瞬时峰值,限流治刷子,读写分离治读压力,分库分表治数据量和写入天花板。五张牌按需组合,但每张都有代价——架构的本质是权衡。


第 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 是最早的分布式事务协议,把提交拆成两个阶段:

  1. Prepare(准备阶段):协调者向所有参与者发「prepare」,每个参与者执行本地事务但不提交,返回「可以提交」或「不能提交」;
  2. Commit /
    Rollback(提交阶段)
    :所有参与者都「可以提交」时,协调者发「commit」让所有人提交;任一参与者「不能提交」,协调者发「rollback」让所有人回滚。

Spring Boot 里可以用 JTA(如 Atomikos、Narayana)实现
2PC,但生产上很少用,原因有二:

  • 同步阻塞:准备阶段所有资源被锁住,直到提交完成,吞吐极低;
  • 协调者单点:协调者宕机会让参与者「卡在
    prepare」状态,需要额外机制恢复。

所以 2PC
只适合「同一系统内跨多个同类型数据库、强一致必须、并发量不大」的场景(如财务系统跨两个
MySQL 库),不适合高并发电商。

5.7
事务消息:本地消息表的「进阶版」

本地消息表有个小缺陷:业务和消息虽然同库同事务,但消息的「投递」依赖后台定时任务,有秒级延迟。RocketMQ
事务消息把「业务 + 消息」的原子性从应用层下沉到 MQ
层,投递更及时:

  1. 生产者发送「半消息」(消费者暂时收不到);
  2. 执行本地事务(下单);
  3. 本地事务成功,提交半消息(消费者可收到);失败则回滚半消息(消费者永远收不到);
  4. 若生产者宕机没来得及提交/回滚,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
分层重构不是「推倒重写」,而是「小步快跑」。一个可落地的重构路径:

  1. 先补测试:给现有 OrderService
    的核心用例写测试(阶段 8),建立安全网;
  2. 提取领域规则:把 OrderService
    里「算金额、校验状态流转」的逻辑,逐步搬进 Order
    聚合根的方法里,每搬一步跑一次测试;
  3. 引入仓储接口:把 OrderMapper
    的直接调用,替换成 OrderRepository
    接口(依赖倒置),实现放基础设施层;
  4. 拆应用服务:让 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_IDMyBatis-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 延伸阅读

  1. Redisson
    官方文档
    https://github.com/redisson/redisson/wiki
    —— 分布式锁、看门狗、集群模式的一手资料。
  2. Apache ShardingSphere
    官方文档
    https://shardingsphere.apache.org/ ——
    分库分表、读写分离、数据治理的权威参考。
  3. 《数据密集型应用系统设计》(DDIA):Martin
    Kleppmann 著 ——
    分布式系统(一致性、复制、分区)的奠基之作,架构师必读。
  4. 《实现领域驱动设计》:Vaughn Vernon 著 —— DDD
    从概念到落地(聚合、限界上下文、事件)的经典。
  5. 《凤凰架构》:周志明著 ——
    从单体到微服务、从强一致到最终一致的演进,中文架构读物里最系统的一本。

文档生成日期:2026-08-14

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