Spring cloud

Spring Cloud 微服务架构全解(下):熔断、事务与消息

2026-07-15 #Spring cloud#微服务

Feign Fallback 解决的是”调用失败后怎么办”,但如果商品服务持续故障,每次调用都等超时——系统会被拖死。Sentinel 在故障扩散前就熔断。Seata 保证跨服务操作的一致性。RocketMQ 把同步链变成异步扇出。本文继续拆解微服务的高级特性。


第一章:Sentinel 熔断限流

1.1 雪崩效应

1
2
3
4
5
6
订单服务 → 商品服务(持续故障)
│ 第1次调用:等10秒超时 → Fallback
│ 第2次调用:等10秒超时 → Fallback
│ 线程池被超时请求占满 → 订单服务也挂了

雪崩:一个服务故障导致整个系统不可用

1.2 熔断器三态

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
┌─────────┐
│ CLOSED │ ← 正常状态,请求放行
└────┬────┘
│ 错误率超阈值

┌─────────┐
│ OPEN │ ← 熔断状态,直接返回Fallback
└────┬────┘
│ 等待恢复时间

┌─────────┐
│HALF_OPEN│ ← 半开,放行一个探测请求
└────┬────┘
│ 成功 → CLOSED
│ 失败 → OPEN

1.3 配置

1
2
3
4
<dependency>
<groupId>com.alibaba.cloud</groupId>
<artifactId>spring-cloud-starter-alibaba-sentinel</artifactId>
</dependency>
1
2
3
4
5
6
7
8
9
10
spring:
cloud:
sentinel:
transport:
dashboard: localhost:8088
port: 8719
eager: true
openfeign:
sentinel:
enabled: true # ★ 开启Feign对Sentinel的支持

1.4 三重保护链路

1
2
3
4
5
6
7
8
9
10
11
请求到达

▼ Sentinel检查
│ OPEN → 直接Fallback(不发起调用)
│ CLOSED → 继续
│ HALF_OPEN → 放行1个探测

▼ Feign调用
│ 超时 → Fallback
│ 5xx → Fallback
│ 200 → 返回结果

1.5 滑动窗口算法

1
2
3
4
5
6
7
8
时间轴 →
┌────┬────┬────┬────┬────┐
│ W1 │ W2 │ W3 │ W4 │ W5 │ 每个窗口100ms
│ 3次 │ 2次 │ 5次 │ 1次 │ 2次 │
└────┴────┴────┴────┴────┘
↑ ↑
└── 最近1秒窗口 ──────┘
总请求数 = 2+5+1+2 = 10

固定窗口在窗口切换瞬间会丢失前一个窗口的数据;滑动窗口始终保留最近 N 个窗口的完整统计。


第二章:Seata AT 分布式事务

2.1 分布式事务的难题

1
2
3
4
5
6
7
订单服务                        商品服务
│ 1.扣减库存(Feign) ──────────>│ UPDATE product SET stock=stock-1
│ <── 成功 ──────────────────│
│ 2.创建订单(本地) │
│ INSERT INTO orders ... ✗失败 │
│ │
│ 库存已扣,订单没创建 → 超卖! │

2.2 Seata 三大角色

角色 职责 示例
TC 全局事务协调者 Seata Server(独立部署)
TM 事务发起者 订单服务(标注@GlobalTransactional
RM 资源管理者 商品服务(执行SQL,记录undo_log)

2.3 AT 模式执行流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
TM: 订单服务
│ ① 开启全局事务,获取XID
│ ② Feign调用商品服务(XID通过Header传递)
│──────────────────────────────────────> RM: 商品服务
│ │ ③ 注册分支事务到TC
│ │ ④ 执行SQL前:记录before image
│ │ ⑤ 执行UPDATE(扣库存)
│ │ ⑥ 执行SQL后:记录after image
│ │ ⑦ 生成undo_log,注册到TC
│ <────────────────────────────────────│
│ ⑧ 本地INSERT(创建订单)✗ 失败
│ ⑨ TM通知TC回滚
│──────────────────────────────────────> TC通知所有RM回滚
│ │ RM用undo_log反向回滚库存
│ │ UPDATE product SET stock=stock+1

2.4 创建订单(全局事务)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
@Override
@GlobalTransactional(name = "cgb-create-order", rollbackFor = Exception.class)
public void createOrder(OrdersEntity entity) {
// ① 远程调用商品服务扣减库存(RM端)
var stockResult = feignProductService.decreaseStock(
entity.getProductId(), entity.getQuantity());
if (stockResult.getCode() != 0) {
throw new EIException("库存扣减失败: " + stockResult.getMsg());
}

// ② 计算总价
entity.setTotalPrice(entity.getUnitPrice()
.multiply(BigDecimal.valueOf(entity.getQuantity())));

// ③ 保存订单(本地TM端)
ordersDao.insert(entity);

// ④ 发送MQ消息
sendOrderStatusMessage(entity, MQTopics.TAG_ORDER_CREATED);
}

如果 ③ 失败:Seata TC 通知商品服务 RM,用 undo_log 自动回滚库存。

2.5 四处全局事务

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// 参团:人数+1 + 扣库存 + 发MQ
@GlobalTransactional(name = "cgb-join-groupbuy")
public void joinGroupBuy(Long groupBuyId, Long userId, Integer quantity) { ... }

// 创建订单:扣库存 + 创建订单 + 发MQ
@GlobalTransactional(name = "cgb-create-order")
public void createOrder(OrdersEntity entity) { ... }

// 取消订单:更新状态 + 回补库存 + 发MQ
@GlobalTransactional(name = "cgb-cancel-order")
public void cancel(String orderId, Long userId) { ... }

// 购物车结算:批量创建订单 + 清空购物车
@GlobalTransactional(name = "cgb-cart-checkout")
public List<OrderVO> checkout(Long userId) { ... }

2.6 RM 端:商品服务扣减库存

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
@Override
public R<?> decreaseStock(Long id, Integer quantity) {
String key = STOCK_KEY_PREFIX + id;

// ① Redis预扣减(高性能防超卖)
Long remain = redisTemplate.opsForValue().decrement(key, quantity);
if (remain != null && remain < 0) {
redisTemplate.opsForValue().increment(key, quantity); // 回滚Redis
return R.fail("库存不足");
}

// ② 数据库扣减(Seata AT自动管理undo_log)
int rows = shangpinDao.decreaseStock(id, quantity);
if (rows == 0) {
redisTemplate.opsForValue().increment(key, quantity); // 回滚Redis
return R.fail("库存不足");
}
return R.ok();
}

Redis + DB 双重扣减:Redis 先扣(高性能防超卖),DB 再扣(持久化)。Seata 回滚 DB 但不会回滚 Redis——需要通过 RocketMQ 消息补偿。

2.7 undo_log 表

1
2
3
4
5
6
7
CREATE TABLE `undo_log` (
`branch_id` BIGINT NOT NULL,
`xid` VARCHAR(128) NOT NULL COMMENT '全局事务ID',
`rollback_info` LONGBLOB NOT NULL COMMENT 'before/after image',
`log_status` INT NOT NULL,
PRIMARY KEY (`branch_id`)
);

回滚时:UPDATE product SET stock = 100 WHERE id = 1(用 beforeImage 反向操作)。


第三章:RocketMQ 异步消息

3.1 同步 vs 异步

1
2
3
4
5
6
7
8
9
10
11
同步链路(慢):
用户支付 → 加积分(200ms) → 更新销量(150ms) → 生成公告(100ms) → 返回
总耗时:450ms(用户等待)

异步扇出(快):
用户支付 → 发送MQ消息(5ms) → 返回成功

├── 用户服务消费 → 加积分(异步)
├── 商品服务消费 → 更新销量(异步)
└── 内容服务消费 → 生成公告(异步)
用户感知耗时:5ms

3.2 消息主题与标签

1
2
3
4
5
6
7
8
9
10
11
12
public interface MQTopics {
String ORDER_STATUS_CHANGE = "ORDER_STATUS_CHANGE";
String GROUPBUY_STATUS_CHANGE = "GROUPBUY_STATUS_CHANGE";

String TAG_ORDER_CREATED = "ORDER_CREATED";
String TAG_ORDER_PAID = "ORDER_PAID";
String TAG_ORDER_CANCELLED = "ORDER_CANCELLED";

String TAG_GROUPBUY_JOINED = "GROUPBUY_JOINED";
String TAG_GROUPBUY_COMPLETED = "GROUPBUY_COMPLETED";
String TAG_GROUPBUY_EXPIRED = "GROUPBUY_EXPIRED";
}

3.3 消息生产者

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
private void sendOrderStatusMessage(OrdersEntity order, String tag) {
try {
OrderStatusMessage msg = new OrderStatusMessage();
msg.setOrderId(order.getOrderNo());
msg.setUserId(order.getUserId());
msg.setProductId(order.getProductId());
msg.setTotalPrice(order.getTotalPrice());
msg.setStatus(order.getStatus());

// destination格式:topic:tag
String destination = MQTopics.ORDER_STATUS_CHANGE + ":" + tag;
rocketMQTemplate.syncSend(destination,
MessageBuilder.withPayload(msg).build());
} catch (Exception e) {
log.error("消息发送失败: orderId={}, tag={}", order.getOrderNo(), tag, e);
}
}

3.4 四个消费者

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
// 消费者1:用户服务消费支付消息(加积分)
@RocketMQMessageListener(
topic = MQTopics.ORDER_STATUS_CHANGE,
selectorExpression = MQTopics.TAG_ORDER_PAID,
consumerGroup = "cgb-user-order-consumer-group"
)
public class UserOrderMessageConsumer implements RocketMQListener<OrderStatusMessage> {
@Override
public void onMessage(OrderStatusMessage message) {
if (message.getStatus() == 1) {
yonghuService.addPoints(message.getUserId(),
message.getTotalPrice().doubleValue());
}
}
}

// 消费者2:商品服务消费订单消息(更新销量)
@RocketMQMessageListener(
topic = MQTopics.ORDER_STATUS_CHANGE,
selectorExpression = "*",
consumerGroup = "cgb-product-order-consumer-group"
)
public class ProductOrderMessageConsumer implements RocketMQListener<OrderStatusMessage> { ... }

// 消费者3:内容服务消费成团消息(生成公告)
@RocketMQMessageListener(
topic = MQTopics.GROUPBUY_STATUS_CHANGE,
selectorExpression = MQTopics.TAG_GROUPBUY_COMPLETED,
consumerGroup = "cgb-content-groupbuy-consumer-group"
)
public class ContentGroupBuyConsumer implements RocketMQListener<GroupBuyMessage> {
@Override
public void onMessage(GroupBuyMessage message) {
NewsEntity news = new NewsEntity();
news.setTitle("🎉 团购成团通知");
news.setContent(String.format("团购ID=%d,共%d人参团成功",
message.getGroupBuyId(), message.getCurrentMemberCount()));
newsService.save(news);
}
}

// 消费者4:团购服务消费过期消息(回补库存)
@RocketMQMessageListener(
topic = MQTopics.GROUPBUY_STATUS_CHANGE,
selectorExpression = "*",
consumerGroup = "cgb-groupbuy-status-consumer-group"
)
public class GroupBuyStatusConsumer implements RocketMQListener<GroupBuyMessage> {
@Override
public void onMessage(GroupBuyMessage message) {
if (message.getStatus() == 2) { // 过期
feignProductService.increaseStock(
message.getProductId(), message.getQuantity());
}
}
}

3.5 消费幂等

RocketMQ 可能重复投递消息,消费者必须幂等:

1
2
3
4
5
6
7
8
9
10
11
@Override
public void onMessage(OrderStatusMessage message) {
String processedKey = "mq:processed:" + message.getOrderId();
// SETNX:第一次设置返回true,重复设置返回false
if (redisTemplate.opsForValue().setIfAbsent(
processedKey, "1", 24, TimeUnit.HOURS)) {
yonghuService.addPoints(message.getUserId(), points);
} else {
log.info("消息已处理,跳过: orderId={}", message.getOrderId());
}
}

3.6 消息流转全景

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
订单服务                      用户服务
│ 支付成功 │
│ │ │
│ ▼ 发送MQ │
│ ORDER_STATUS_CHANGE:ORDER_PAID
│ │ │
│ │ ┌────────────────────────┘
│ ▼ ▼ 消费者1 → 加积分
│ 用户服务.onMessage()

│ │ ┌────────────────────────┐
│ ▼ ▼ 消费者2 → 更新销量
│ 商品服务.onMessage()

│ 发送MQ
│ GROUPBUY_STATUS_CHANGE:GROUPBUY_COMPLETED
│ │
│ ▼ 消费者3 → 生成公告
│ 内容服务.onMessage()

第四章:雪花算法 ID 生成

4.1 分布式唯一 ID

在 Spring Cloud 2024 Demo 项目中,订单号用雪花算法生成:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
import cn.hutool.core.util.IdUtil;

@Override
@Transactional(rollbackFor = Exception.class)
public R<Order> createOrder(Long userId, Long productId, Integer quantity) {
// 远程调用商品服务查询和扣减
R<?> productResult = productFeignClient.getProduct(productId);
R<Void> stockResult = productFeignClient.deductStock(productId, quantity);

// ★ 雪花算法生成分布式唯一订单号
Order order = new Order();
order.setOrderNo(IdUtil.getSnowflakeNextIdStr());
order.setUserId(userId);
order.setProductId(productId);
save(order);
return R.ok(order);
}

4.2 雪花算法原理

1
2
3
4
┌─────────────────────────────────────────────────┐
│ 1位符号 │ 41位时间戳 │ 10位机器ID │ 12位序列号 │
│ (不用) │ (约69年) │ (1024台) │ (4096/ms) │
└─────────────────────────────────────────────────┘
  • 时间戳:保证趋势递增(数据库索引友好)
  • 机器ID:分布式唯一(每台机器不同)
  • 序列号:同一毫秒内的递增序号

第五章:乐观锁防超卖

5.1 MyBatis Plus 乐观锁

1
2
3
4
5
6
7
8
9
10
@Configuration
public class MyBatisPlusConfig {
@Bean
public MybatisPlusInterceptor mybatisPlusInterceptor() {
MybatisPlusInterceptor interceptor = new MybatisPlusInterceptor();
interceptor.addInnerInterceptor(new PaginationInnerInterceptor(DbType.MYSQL));
interceptor.addInnerInterceptor(new OptimisticLockerInnerInterceptor());
return interceptor;
}
}

5.2 Product 实体

1
2
3
4
5
6
7
8
9
10
@Data
@TableName("product")
public class Product extends BaseEntity {
private Long id;
private String name;
private BigDecimal price;
private Integer stock;
/** 版本号(乐观锁) */
private Integer version;
}

5.3 扣减库存

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
@Override
@Transactional(rollbackFor = Exception.class)
public R<Void> deductStock(Long productId, Integer quantity) {
Product product = getById(productId);
if (product.getStock() < quantity) {
return R.failed(ResultCode.STOCK_NOT_ENOUGH);
}

product.setStock(product.getStock() - quantity);
// ★ 乐观锁:updateById 会自动加 WHERE version = ?
boolean success = updateById(product);

if (!success) {
return R.failed("库存扣减失败,请重试");
}

redisTemplate.delete("product:info:" + productId); // 删除缓存
return R.ok();
}

乐观锁 vs 悲观锁:乐观锁不阻塞,CAS 重试;悲观锁 SELECT FOR UPDATE 阻塞。


第六章:Token 黑名单

6.1 登出时加入黑名单

1
2
3
4
5
6
7
8
9
10
@PostMapping("/logout")
public R<?> logout(HttpServletRequest request) {
String token = JwtUtil.getTokenFromRequest(request);
if (token != null && JwtUtil.isTokenValid(token)) {
// 将Token加入Redis黑名单(剩余有效期内)
long expire = JwtUtil.getRemainingTime(token);
redisUtil.set("blacklist:" + token, "true", expire, TimeUnit.MILLISECONDS);
}
return R.ok("登出成功");
}

6.2 网关检查黑名单

1
2
3
4
5
// 网关过滤器中检查
Boolean isBlacklisted = redisTemplate.hasKey("token:blacklist:" + token);
if (Boolean.TRUE.equals(isBlacklisted)) {
return unauthorized(exchange, "Token已失效");
}

总结

技术 解决的问题 核心实现
Sentinel 雪崩效应 滑动窗口 + 熔断三态
Seata AT 分布式事务一致性 @GlobalTransactional + undo_log
RocketMQ 异步解耦 topic:tag + 消费者组
雪花算法 分布式唯一ID 时间戳 + 机器ID + 序列号
乐观锁 并发扣减 version字段 + CAS
Token黑名单 登出失效 Redis + TTL

三道防线:Sentinel 熔断(最快切断)→ Feign 超时(防止无限等待)→ Fallback 降级(返回友好响应)。Seata 保证数据一致性,RocketMQ 解耦异步逻辑。这就是微服务高级特性的完整图景。

评论
分享