Spring Boot高并发架构实战:接口限流与分布式事务设计方案

Spring Boot框架在高并发场景下面临接口限流和分布式事务两大核心挑战。微服务架构拆分后,跨服务数据一致性问题和突发流量冲击成为系统稳定性的主要风险点。本文从实际电商系统出发,展示Redis令牌桶限流实现、Seata分布式事务集成和Saga补偿模式的设计方案。

接口限流实现:Redis令牌桶算法与AOP集成

高并发系统中接口限流是第一道防线。基于Redis + Lua脚本实现的令牌桶算法,具备原子性和高性能,适用于分布式环境:

-- Redis令牌桶限流 Lua脚本
-- resources/lua/token_bucket.lua

local key = KEYS[1]
local capacity = tonumber(ARGV[1])
local rate = tonumber(ARGV[2])
local now = tonumber(ARGV[3])
local requested = tonumber(ARGV[4])

local bucket = redis.call('hmget', key, 'tokens', 'timestamp')
local tokens = tonumber(bucket[1]) or capacity
local last_time = tonumber(bucket[2]) or now

local delta = math.max(0, now - last_time)
local new_tokens = math.min(capacity, tokens + (delta * rate / 1000))

if new_tokens >= requested then
    new_tokens = new_tokens - requested
    redis.call('hmset', key, 'tokens', new_tokens, 'timestamp', now)
    redis.call('expire', key, 60)
    return 1
else
    redis.call('hmset', key, 'tokens', new_tokens, 'timestamp', now)
    redis.call('expire', key, 60)
    return 0
end

Java侧封装为注解+AOP拦截器,支持方法级限流配置:

// RateLimit注解定义
@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
public @interface RateLimit {
    String key() default "";
    int capacity() default 100;
    double rate() default 10.0;
    int requested() default 1;
    String message() default "请求过于频繁,请稍后重试";
}

// AOP切面实现
@Aspect
@Component
public class RateLimitAspect {

    @Autowired
    private RedisTemplate<String, Object> redisTemplate;

    @Around("@annotation(rateLimit)")
    public Object around(ProceedingJoinPoint joinPoint, RateLimit rateLimit) throws Throwable {
        String key = buildKey(joinPoint, rateLimit);

        DefaultRedisScript<Long> script = new DefaultRedisScript<>();
        script.setScriptSource(new ResourceScriptSource(
            new ClassPathResource("lua/token_bucket.lua")));
        script.setResultType(Long.class);

        long now = System.currentTimeMillis();
        Long result = redisTemplate.execute(
            script,
            Collections.singletonList(key),
            String.valueOf(rateLimit.capacity()),
            String.valueOf(rateLimit.rate()),
            String.valueOf(now),
            String.valueOf(rateLimit.requested())
        );

        if (result == null || result == 0) {
            throw new RateLimitException(rateLimit.message());
        }
        return joinPoint.proceed();
    }

    private String buildKey(ProceedingJoinPoint joinPoint, RateLimit rateLimit) {
        MethodSignature sig = (MethodSignature) joinPoint.getSignature();
        String method = sig.getDeclaringType().getSimpleName()
            + ":" + sig.getMethod().getName();

        if (!rateLimit.key().isEmpty()) {
            EvaluationContext ctx = new StandardEvaluationContext();
            Object[] args = joinPoint.getArgs();
            String[] paramNames = sig.getParameterNames();
            for (int i = 0; i < paramNames.length; i++) {
                ctx.setVariable(paramNames[i], args[i]);
            }
            String parsedKey = new SpelExpressionParser()
                .parseExpression(rateLimit.key())
                .getValue(ctx, String.class);
            return "rate_limit:" + method + ":" + parsedKey;
        }
        return "rate_limit:" + method;
    }
}

// Controller使用示例
@RestController
@RequestMapping("/api/orders")
public class OrderController {

    @PostMapping
    @RateLimit(key = "#userId", capacity = 20, rate = 2.0,
               message = "下单频率超限,每秒最多2次")
    public Result<OrderVO> createOrder(@RequestBody OrderRequest request,
                                         @RequestHeader("X-User-Id") String userId) {
        return Result.success(orderService.create(request, userId));
    }
}

分布式事务方案选型:Seata AT模式集成实践

微服务架构中,订单创建涉及库存扣减、优惠券核销、积分发放等多个服务。传统XA事务性能差且不适用于微服务。Seata的AT(Automatic Transaction)模式通过SQL解析自动生成补偿操作,对业务代码侵入小:

// application.yml配置
// seata:
//   enabled: true
//   application-id: order-service
//   tx-service-group: yunthe_tx_group
//   service:
//     vgroup-mapping:
//       yunthe_tx_group: default
//     grouplist:
//       default: 127.0.0.1:8091

// 业务代码中使用@GlobalTransactional
@Service
public class OrderServiceImpl implements OrderService {

    @Autowired
    private InventoryClient inventoryClient;
    @Autowired
    private CouponClient couponClient;
    @Autowired
    private PointsClient pointsClient;
    @Autowired
    private OrderMapper orderMapper;

    @Override
    @GlobalTransactional(timeoutMills = 60000, name = "create-order")
    public OrderVO create(OrderRequest request, String userId) {
        // 1. 创建订单(本地事务)
        Order order = new Order();
        order.setUserId(userId);
        order.setProductId(request.getProductId());
        order.setQuantity(request.getQuantity());
        order.setAmount(request.getAmount());
        order.setStatus("CREATED");
        orderMapper.insert(order);

        // 2. 扣减库存(远程调用,RM自动注册分支事务)
        inventoryClient.deduct(request.getProductId(), request.getQuantity());

        // 3. 核销优惠券(远程调用)
        if (request.getCouponId() != null) {
            couponClient.redeem(request.getCouponId(), userId, order.getId());
        }

        // 4. 发放积分(远程调用)
        pointsClient.addPoints(userId, order.getAmount().intValue());

        // 任何一步失败,Seata TC自动协调所有分支回滚
        return OrderVO.from(order);
    }
}

Seata AT模式的核心机制:每个分支事务执行时,Seata自动记录数据变更前后的快照(beforeImage / afterImage),存入undo_log表。全局回滚时根据undo_log反向执行SQL恢复数据。需要注意,AT模式要求所有参与方使用支持ACID的数据库,且每张业务表必须有主键。

Saga补偿模式:长事务的替代方案

当事务涉及第三方系统或耗时操作时,Seata AT的全局锁会导致性能瓶颈。Saga模式通过补偿事务替代回滚,更适合长链路业务:

// Saga编排式实现 - 订单创建Saga
@Component
public class OrderCreateSaga {

    @Autowired
    private InventoryClient inventoryClient;
    @Autowired
    private CouponClient couponClient;
    @Autowired
    private PointsClient pointsClient;
    @Autowired
    private OrderMapper orderMapper;

    private final SagaOrchestrator sagaOrchestrator;

    public OrderCreateSaga(SagaOrchestrator sagaOrchestrator) {
        this.sagaOrchestrator = sagaOrchestrator;
    }

    public OrderVO execute(OrderRequest request, String userId) {
        String sagaId = UUID.randomUUID().toString();

        List<SagaStep> steps = Arrays.asList(
            SagaStep.builder()
                .name("create-order")
                .action(() -> {
                    Order order = buildOrder(request, userId);
                    orderMapper.insert(order);
                    return order;
                })
                .compensation((result) -> {
                    Order order = (Order) result;
                    orderMapper.deleteById(order.getId());
                })
                .build(),

            SagaStep.builder()
                .name("deduct-inventory")
                .action(() -> inventoryClient.deduct(
                    request.getProductId(), request.getQuantity()))
                .compensation((result) -> inventoryClient.restore(
                    request.getProductId(), request.getQuantity()))
                .build(),

            SagaStep.builder()
                .name("redeem-coupon")
                .action(() -> {
                    if (request.getCouponId() == null) return null;
                    return couponClient.redeem(
                        request.getCouponId(), userId, sagaId);
                })
                .compensation((result) -> {
                    if (result != null) {
                        couponClient.cancelRedeem(request.getCouponId(), sagaId);
                    }
                })
                .build(),

            SagaStep.builder()
                .name("add-points")
                .action(() -> pointsClient.addPoints(
                    userId, request.getAmount().intValue()))
                .compensation((result) -> pointsClient.deductPoints(
                    userId, request.getAmount().intValue()))
                .build()
        );

        SagaResult result = sagaOrchestrator.execute(sagaId, steps);

        if (result.isSuccess()) {
            return OrderVO.from((Order) result.getStepResult("create-order"));
        } else {
            throw new BusinessException("订单创建失败: " + result.getErrorMessage());
        }
    }
}

// Saga编排器核心逻辑
@Component
public class SagaOrchestrator {

    public SagaResult execute(String sagaId, List<SagaStep> steps) {
        Map<String, Object> stepResults = new HashMap<>();
        List<Integer> completedSteps = new ArrayList<>();

        try {
            for (int i = 0; i < steps.size(); i++) {
                SagaStep step = steps.get(i);
                Object result = step.getAction().get();
                stepResults.put(step.getName(), result);
                completedSteps.add(i);
            }
            return SagaResult.success(stepResults);

        } catch (Exception e) {
            // 反向执行补偿
            for (int i = completedSteps.size() - 1; i >= 0; i--) {
                int stepIndex = completedSteps.get(i);
                SagaStep step = steps.get(stepIndex);
                try {
                    step.getCompensation().accept(
                        stepResults.get(step.getName()));
                } catch (Exception compEx) {
                    log.error("Saga补偿失败: sagaId={}, step={}",
                        sagaId, step.getName(), compEx);
                }
            }
            return SagaResult.failure(e.getMessage(), stepResults);
        }
    }
}

RocketMQ事务消息:异步场景的最终一致性保障

对于异步场景(如支付成功后通知物流),使用消息中间件的事务消息机制保障最终一致性。RocketMQ事务消息通过半消息+本地事务+回查机制确保消息与本地事务的原子性:

@Service
public class PaymentEventListener {

    @Autowired
    private TransactionMQProducer producer;

    public void sendPaymentSuccess(String orderId, BigDecimal amount)
            throws MQClientException {
        Message msg = new Message(
            "PAYMENT_TOPIC", "PAYMENT_SUCCESS", orderId,
            JSON.toJSONString(Map.of("orderId", orderId, "amount", amount))
                .getBytes(StandardCharsets.UTF_8)
        );
        producer.sendMessageInTransaction(msg, orderId);
    }

    @Component
    public class LocalTransactionExecutor implements TransactionListener {

        @Autowired
        private PaymentMapper paymentMapper;

        @Override
        public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
            String orderId = (String) arg;
            try {
                paymentMapper.updateStatus(orderId, "PAID");
                return LocalTransactionState.COMMIT_MESSAGE;
            } catch (Exception e) {
                return LocalTransactionState.ROLLBACK_MESSAGE;
            }
        }

        @Override
        public LocalTransactionState checkLocalTransaction(MessageExt msg) {
            String orderId = msg.getKeys();
            Payment payment = paymentMapper.selectByOrderId(orderId);
            if (payment != null && "PAID".equals(payment.getStatus())) {
                return LocalTransactionState.COMMIT_MESSAGE;
            } else if (payment != null && "FAILED".equals(payment.getStatus())) {
                return LocalTransactionState.ROLLBACK_MESSAGE;
            }
            return LocalTransactionState.UNKNOW;
        }
    }
}

// 消费端(物流服务)- 幂等处理
@RocketMQMessageListener(
    topic = "PAYMENT_TOPIC",
    consumerGroup = "logistics-consumer-group",
    selectorExpression = "PAYMENT_SUCCESS"
)
@Component
public class LogisticsConsumer implements RocketMQListener<String> {

    @Override
    public void onMessage(String message) {
        JSONObject data = JSON.parseObject(message);
        String orderId = data.getString("orderId");

        // 幂等性检查
        if (logisticsService.existsByOrderId(orderId)) {
            return;
        }
        logisticsService.createShipment(orderId);
    }
}

事务消息方案将分布式事务拆解为本地事务+消息可靠投递,牺牲强一致性换取系统吞吐量。消费端通过幂等设计处理重复消息,最终达到数据一致。在高并发场景下,这种方案的性能远优于同步的分布式事务方案。服务治理框架通过拦截器自动注入消息发送逻辑,实现业务代码与基础设施的解耦。业务中台建设中,事务消息的标准封装可以作为基础设施层能力,供上层业务方直接复用。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot-gao-bing-fa-jia-gou-shi-zhan-jie-kou-xian-liu-yu/

(0)
小编小编
上一篇 9小时前
下一篇 9小时前

相关推荐