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
|
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) { 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())));
ordersDao.insert(entity);
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
| @GlobalTransactional(name = "cgb-join-groupbuy") public void joinGroupBuy(Long groupBuyId, Long userId, Integer quantity) { ... }
@GlobalTransactional(name = "cgb-create-order") public void createOrder(OrdersEntity entity) { ... }
@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; Long remain = redisTemplate.opsForValue().decrement(key, quantity); if (remain != null && remain < 0) { redisTemplate.opsForValue().increment(key, quantity); return R.fail("库存不足"); } int rows = shangpinDao.decreaseStock(id, quantity); if (rows == 0) { redisTemplate.opsForValue().increment(key, quantity); 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());
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
| @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()); } } }
@RocketMQMessageListener( topic = MQTopics.ORDER_STATUS_CHANGE, selectorExpression = "*", consumerGroup = "cgb-product-order-consumer-group" ) public class ProductOrderMessageConsumer implements RocketMQListener<OrderStatusMessage> { ... }
@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); } }
@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(); 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); 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)) { 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 解耦异步逻辑。这就是微服务高级特性的完整图景。