Spring Boot微服务高并发架构设计:从线程池到分布式限流的实战方案

高并发架构的本质是资源管控

Spring Boot微服务在高并发场景下的核心矛盾,不是”能扛多少QPS”,而是”在资源有限的前提下,如何保证核心链路不被非核心流量打垮”。从线程池隔离、接口限流、异步削峰到分布式限流,每一层都在解决同一个问题:让有限的计算资源优先服务最高优先级的请求。

线程池隔离:不同业务不同池

Spring Boot内嵌Tomcat默认200个工作线程,所有接口共享。当某个慢接口占满线程池后,其他接口全部排队等待。Hystrix的线程池隔离思想在Sentinel中同样适用:

@Configuration
public class ThreadPoolConfig {

    @Bean("orderExecutor")
    public ThreadPoolExecutor orderExecutor() {
        return new ThreadPoolExecutor(
            10,                              // 核心线程数
            50,                              // 最大线程数
            60, TimeUnit.SECONDS,            // 空闲存活时间
            new LinkedBlockingQueue<>(200),   // 队列容量
            new CustomNamedThreadFactory("order-pool"),
            new ThreadPoolExecutor.CallerRunsPolicy()  // 拒绝策略
        );
    }

    @Bean("reportExecutor")
    public ThreadPoolExecutor reportExecutor() {
        // 报表生成类任务:低优先级,小池子
        return new ThreadPoolExecutor(
            2, 10, 120, TimeUnit.SECONDS,
            new LinkedBlockingQueue<>(50),
            new CustomNamedThreadFactory("report-pool"),
            new ThreadPoolExecutor.DiscardOldestPolicy()
        );
    }
}

// 自定义ThreadFactory
public class CustomNamedThreadFactory implements ThreadFactory {
    private final AtomicInteger counter = new AtomicInteger(1);
    private final String prefix;

    public CustomNamedThreadFactory(String prefix) {
        this.prefix = prefix;
    }

    @Override
    public Thread newThread(Runnable r) {
        Thread t = new Thread(r, prefix + "-" + counter.getAndIncrement());
        t.setDaemon(false);
        return t;
    }
}

线程池参数的确定方法:

CPU密集型任务:核心线程数 = CPU核数 + 1
IO密集型任务:核心线程数 = CPU核数 x 2 / (1 – 阻塞系数)

实际生产中建议用压力测试工具(JMeter/wrk)逐步加压,观察线程池活跃数和队列堆积情况来微调。

接口级限流:Sentinel + 注解驱动

// 自定义限流注解
@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
public @interface RateLimit {
    String resource();           // 资源名
    int qps() default 100;      // QPS阈值
    String fallback() default ""; // 降级方法
}

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

    @Around("@annotation(rateLimit)")
    public Object around(ProceedingJoinPoint pjp, RateLimit rateLimit) throws Throwable {
        Entry entry = null;
        try {
            entry = SphU.entry(rateLimit.resource());
            return pjp.proceed();
        } catch (BlockException e) {
            if (!rateLimit.fallback().isEmpty()) {
                return executeFallback(pjp, rateLimit.fallback());
            }
            throw new RateLimitException("接口限流: " + rateLimit.resource());
        } finally {
            if (entry != null) {
                entry.exit();
            }
        }
    }

    private Object executeFallback(ProceedingJoinPoint pjp, String methodName)
            throws Throwable {
        MethodSignature signature = (MethodSignature) pjp.getSignature();
        Method method = Arrays.stream(pjp.getTarget().getClass().getMethods())
            .filter(m -> m.getName().equals(methodName))
            .findFirst()
            .orElseThrow();
        return method.invoke(pjp.getTarget(), pjp.getArgs());
    }
}

// 控制器使用
@RestController
@RequestMapping("/api/orders")
public class OrderController {

    @PostMapping("/create")
    @RateLimit(resource = "order-create", qps = 500, fallback = "createFallback")
    public Result<Order> createOrder(@RequestBody OrderRequest req) {
        return Result.ok(orderService.create(req));
    }

    public Result<Order> createFallback(OrderRequest req) {
        return Result.fail("系统繁忙,请稍后重试");
    }
}

异步削峰:RabbitMQ + 死信队列

秒杀场景下同步下单扛不住瞬时流量,需要用消息队列做削峰填谷:

// 消息队列配置
@Configuration
public class RabbitMQConfig {

    // 下单交换机
    @Bean
    public DirectExchange orderExchange() {
        return new DirectExchange("order.exchange", true, false);
    }

    // 主队列
    @Bean
    public Queue orderQueue() {
        return QueueBuilder.durable("order.queue")
            .withArgument("x-dead-letter-exchange", "order.dlx")
            .withArgument("x-dead-letter-routing-key", "order.failed")
            .withArgument("x-max-length", 10000)  // 队列最大长度
            .withArgument("x-message-ttl", 30000)  // 消息TTL 30秒
            .build();
    }

    // 死信队列
    @Bean
    public Queue orderDeadLetterQueue() {
        return QueueBuilder.durable("order.dlq").build();
    }

    @Bean
    public Binding orderBinding() {
        return BindingBuilder.bind(orderQueue())
            .to(orderExchange()).with("order.create");
    }
}

// 生产者:异步发送下单消息
@Service
public class OrderProducer {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    public void asyncCreateOrder(OrderRequest req) {
        CorrelationData correlation = new CorrelationData(
            UUID.randomUUID().toString()
        );
        rabbitTemplate.convertAndSend(
            "order.exchange", "order.create", req, correlation
        );
    }
}

// 消费者:限速消费
@Component
public class OrderConsumer {

    @Autowired
    private OrderService orderService;

    @RabbitListener(queues = "order.queue", concurrency = "5-10")
    public void handleOrder(OrderRequest req, Channel channel,
                           @Header(AmqpHeaders.DELIVERY_TAG) long tag) {
        try {
            orderService.create(req);
            channel.basicAck(tag, false);
        } catch (Exception e) {
            channel.basicNack(tag, false, false);  // 拒绝进入死信
        }
    }
}

concurrency = “5-10″表示最少5个消费者线程,最大10个。RabbitMQ的prefetch默认250,在高并发场景建议设为20-50,避免消费者过载。

分布式限流:Redis + Lua脚本

单机限流在多实例部署下无效,需要基于Redis实现集群级限流:

@Component
public class DistributedRateLimiter {

    @Autowired
    private StringRedisTemplate redisTemplate;

    // 滑动窗口限流Lua脚本
    private static final String LIMIT_SCRIPT =
        "local key = KEYS[1] " +
        "local limit = tonumber(ARGV[1]) " +
        "local window = tonumber(ARGV[2]) " +
        "local now = tonumber(ARGV[3]) " +
        "local windowStart = now - window " +
        "redis.call('ZREMRANGEBYSCORE', key, 0, windowStart) " +
        "local count = redis.call('ZCARD', key) " +
        "if count < limit then " +
        "  redis.call('ZADD', key, now, now .. ':' .. math.random()) " +
        "  redis.call('EXPIRE', key, window / 1000) " +
        "  return 1 " +
        "else " +
        "  return 0 " +
        "end";

    public boolean tryAcquire(String resource, int qps, int windowSeconds) {
        String key = "rate_limit:" + resource;
        long now = System.currentTimeMillis();
        Long result = redisTemplate.execute(
            new DefaultRedisScript<>(LIMIT_SCRIPT, Long.class),
            Collections.singletonList(key),
            String.valueOf(qps),
            String.valueOf(windowSeconds * 1000),
            String.valueOf(now)
        );
        return result != null && result == 1L;
    }
}

优雅停机与流量预热

Spring Boot应用重启时,正在处理的请求不能直接中断:

# application.yml
server:
  shutdown: graceful    # 优雅停机

spring:
  lifecycle:
    timeout-per-shutdown-phase: 30s  # 最大等待时间

# Kubernetes配置
# spec.template.spec.containers.lifecycle.preStop:
#   exec:
#     command: ["/bin/sh", "-c", "sleep 15"]  # 等待Service摘除

流量预热:新启动的实例先用小流量跑热JVM,逐步放开:

// 基于Sentinel的流量预热
FlowRule warmUpRule = new FlowRule();
warmUpRule.setResource("order-api");
warmUpRule.setCount(500);       // 预热完成后QPS
warmUpRule.setGrade(RuleConstant.FLOW_GRADE_QPS);
warmUpRule.setWarmUpPeriodSec(60);  // 预热时长60秒
FlowRuleManager.loadRules(Collections.singletonList(warmUpRule));

高并发架构设计没有银弹,每一层防护都有性能开销。关键是根据业务场景在吞吐量和稳定性之间找到平衡点——宁可让部分请求快速失败,也不能让整个系统雪崩。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot-wei-fu-wu-gao-bing-fa-jia-gou-she-ji-cong-xian/

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

相关推荐

Spring Boot微服务高并发架构设计:从线程池到分布式限流的实战方案

高并发架构的本质是资源管控

Spring Boot微服务在高并发场景下的核心矛盾,不是”能扛多少QPS”,而是”在资源有限的前提下,如何保证核心链路不被非核心流量打垮”。从线程池隔离、接口限流、异步削峰到分布式限流,每一层都在解决同一个问题:让有限的计算资源优先服务最高优先级的请求。

线程池隔离:不同业务不同池

Spring Boot内嵌Tomcat默认200个工作线程,所有接口共享。当某个慢接口占满线程池后,其他接口全部排队等待。Hystrix的线程池隔离思想在Sentinel中同样适用:

@Configuration
public class ThreadPoolConfig {

    @Bean("orderExecutor")
    public ThreadPoolExecutor orderExecutor() {
        return new ThreadPoolExecutor(
            10,                              // 核心线程数
            50,                              // 最大线程数
            60, TimeUnit.SECONDS,            // 空闲存活时间
            new LinkedBlockingQueue<>(200),   // 队列容量
            new CustomNamedThreadFactory("order-pool"),
            new ThreadPoolExecutor.CallerRunsPolicy()  // 拒绝策略
        );
    }

    @Bean("reportExecutor")
    public ThreadPoolExecutor reportExecutor() {
        // 报表生成类任务:低优先级,小池子
        return new ThreadPoolExecutor(
            2, 10, 120, TimeUnit.SECONDS,
            new LinkedBlockingQueue<>(50),
            new CustomNamedThreadFactory("report-pool"),
            new ThreadPoolExecutor.DiscardOldestPolicy()
        );
    }
}

// 自定义ThreadFactory
public class CustomNamedThreadFactory implements ThreadFactory {
    private final AtomicInteger counter = new AtomicInteger(1);
    private final String prefix;

    public CustomNamedThreadFactory(String prefix) {
        this.prefix = prefix;
    }

    @Override
    public Thread newThread(Runnable r) {
        Thread t = new Thread(r, prefix + "-" + counter.getAndIncrement());
        t.setDaemon(false);
        return t;
    }
}

线程池参数的确定方法:

CPU密集型任务:核心线程数 = CPU核数 + 1
IO密集型任务:核心线程数 = CPU核数 x 2 / (1 – 阻塞系数)

实际生产中建议用压力测试工具(JMeter/wrk)逐步加压,观察线程池活跃数和队列堆积情况来微调。

接口级限流:Sentinel + 注解驱动

// 自定义限流注解
@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
public @interface RateLimit {
    String resource();           // 资源名
    int qps() default 100;      // QPS阈值
    String fallback() default ""; // 降级方法
}

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

    @Around("@annotation(rateLimit)")
    public Object around(ProceedingJoinPoint pjp, RateLimit rateLimit) throws Throwable {
        Entry entry = null;
        try {
            entry = SphU.entry(rateLimit.resource());
            return pjp.proceed();
        } catch (BlockException e) {
            if (!rateLimit.fallback().isEmpty()) {
                return executeFallback(pjp, rateLimit.fallback());
            }
            throw new RateLimitException("接口限流: " + rateLimit.resource());
        } finally {
            if (entry != null) {
                entry.exit();
            }
        }
    }

    private Object executeFallback(ProceedingJoinPoint pjp, String methodName)
            throws Throwable {
        MethodSignature signature = (MethodSignature) pjp.getSignature();
        Method method = Arrays.stream(pjp.getTarget().getClass().getMethods())
            .filter(m -> m.getName().equals(methodName))
            .findFirst()
            .orElseThrow();
        return method.invoke(pjp.getTarget(), pjp.getArgs());
    }
}

// 控制器使用
@RestController
@RequestMapping("/api/orders")
public class OrderController {

    @PostMapping("/create")
    @RateLimit(resource = "order-create", qps = 500, fallback = "createFallback")
    public Result<Order> createOrder(@RequestBody OrderRequest req) {
        return Result.ok(orderService.create(req));
    }

    public Result<Order> createFallback(OrderRequest req) {
        return Result.fail("系统繁忙,请稍后重试");
    }
}

异步削峰:RabbitMQ + 死信队列

秒杀场景下同步下单扛不住瞬时流量,需要用消息队列做削峰填谷:

// 消息队列配置
@Configuration
public class RabbitMQConfig {

    // 下单交换机
    @Bean
    public DirectExchange orderExchange() {
        return new DirectExchange("order.exchange", true, false);
    }

    // 主队列
    @Bean
    public Queue orderQueue() {
        return QueueBuilder.durable("order.queue")
            .withArgument("x-dead-letter-exchange", "order.dlx")
            .withArgument("x-dead-letter-routing-key", "order.failed")
            .withArgument("x-max-length", 10000)  // 队列最大长度
            .withArgument("x-message-ttl", 30000)  // 消息TTL 30秒
            .build();
    }

    // 死信队列
    @Bean
    public Queue orderDeadLetterQueue() {
        return QueueBuilder.durable("order.dlq").build();
    }

    @Bean
    public Binding orderBinding() {
        return BindingBuilder.bind(orderQueue())
            .to(orderExchange()).with("order.create");
    }
}

// 生产者:异步发送下单消息
@Service
public class OrderProducer {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    public void asyncCreateOrder(OrderRequest req) {
        CorrelationData correlation = new CorrelationData(
            UUID.randomUUID().toString()
        );
        rabbitTemplate.convertAndSend(
            "order.exchange", "order.create", req, correlation
        );
    }
}

// 消费者:限速消费
@Component
public class OrderConsumer {

    @Autowired
    private OrderService orderService;

    @RabbitListener(queues = "order.queue", concurrency = "5-10")
    public void handleOrder(OrderRequest req, Channel channel,
                           @Header(AmqpHeaders.DELIVERY_TAG) long tag) {
        try {
            orderService.create(req);
            channel.basicAck(tag, false);
        } catch (Exception e) {
            channel.basicNack(tag, false, false);  // 拒绝进入死信
        }
    }
}

concurrency = “5-10″表示最少5个消费者线程,最大10个。RabbitMQ的prefetch默认250,在高并发场景建议设为20-50,避免消费者过载。

分布式限流:Redis + Lua脚本

单机限流在多实例部署下无效,需要基于Redis实现集群级限流:

@Component
public class DistributedRateLimiter {

    @Autowired
    private StringRedisTemplate redisTemplate;

    // 滑动窗口限流Lua脚本
    private static final String LIMIT_SCRIPT =
        "local key = KEYS[1] " +
        "local limit = tonumber(ARGV[1]) " +
        "local window = tonumber(ARGV[2]) " +
        "local now = tonumber(ARGV[3]) " +
        "local windowStart = now - window " +
        "redis.call('ZREMRANGEBYSCORE', key, 0, windowStart) " +
        "local count = redis.call('ZCARD', key) " +
        "if count < limit then " +
        "  redis.call('ZADD', key, now, now .. ':' .. math.random()) " +
        "  redis.call('EXPIRE', key, window / 1000) " +
        "  return 1 " +
        "else " +
        "  return 0 " +
        "end";

    public boolean tryAcquire(String resource, int qps, int windowSeconds) {
        String key = "rate_limit:" + resource;
        long now = System.currentTimeMillis();
        Long result = redisTemplate.execute(
            new DefaultRedisScript<>(LIMIT_SCRIPT, Long.class),
            Collections.singletonList(key),
            String.valueOf(qps),
            String.valueOf(windowSeconds * 1000),
            String.valueOf(now)
        );
        return result != null && result == 1L;
    }
}

优雅停机与流量预热

Spring Boot应用重启时,正在处理的请求不能直接中断:

# application.yml
server:
  shutdown: graceful    # 优雅停机

spring:
  lifecycle:
    timeout-per-shutdown-phase: 30s  # 最大等待时间

# Kubernetes配置
# spec.template.spec.containers.lifecycle.preStop:
#   exec:
#     command: ["/bin/sh", "-c", "sleep 15"]  # 等待Service摘除

流量预热:新启动的实例先用小流量跑热JVM,逐步放开:

// 基于Sentinel的流量预热
FlowRule warmUpRule = new FlowRule();
warmUpRule.setResource("order-api");
warmUpRule.setCount(500);       // 预热完成后QPS
warmUpRule.setGrade(RuleConstant.FLOW_GRADE_QPS);
warmUpRule.setWarmUpPeriodSec(60);  // 预热时长60秒
FlowRuleManager.loadRules(Collections.singletonList(warmUpRule));

高并发架构设计没有银弹,每一层防护都有性能开销。关键是根据业务场景在吞吐量和稳定性之间找到平衡点——宁可让部分请求快速失败,也不能让整个系统雪崩。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot-wei-fu-wu-gao-bing-fa-jia-gou-she-ji-cong-xian/

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

相关推荐