Spring Cloud Gateway动态路由与负载均衡在高并发场景的流量清洗与过载防护

Spring Cloud Gateway动态路由与负载均衡在高并发场景的流量清洗与过载防护

一、概述

Spring Cloud Gateway作为微服务网关,承担着统一入口、路由转发、跨横切关注点的核心职责。在高并发场景下,动态路由的实时变更、负载均衡的策略选择、流量突增时的过载防护,是保障系统稳定性的三大命门。

动态路由的灵活性意味着路由规则可以随时调整,但这也可能引入配置错误、路由风暴等风险。负载均衡需要从简单的轮询升级为感知实例健康状态和实时负载的自适应策略。过载防护则要在流量超过系统承载能力时,优雅地拒绝而非崩溃。

本文聚焦这三者的联动实践,给出高并发场景下的完整流量治理方案。

二、核心原理

2.1 动态路由的三层刷新机制

触发层:  Nacos配置变更 / API调用 / 定时刷新
    ↓
中间层:  RefreshRoutesEvent 事件发布
    ↓
执行层:  CachingRouteLocator 刷新路由缓存

2.2 负载均衡的决策链路

请求到达
    ↓
RoutePredicateHandlerMapping 匹配路由
    ↓
LoadBalancerClientFilter 触发负载均衡
    ↓
ReactorLoadBalancer 选择实例
    ↓
NettyRoutingFilter 转发请求

2.3 过载防护的三个阶段

阶段 策略 响应
事前 容量评估 + 限流阈值设定 配置预热
事中 令牌桶限流 + 熔断降级 429/503
事后 熔断恢复 + 流量回放 渐进恢复

三、实战配置

3.1 基础配置

spring:
  cloud:
    gateway:
      routes:
        - id: order-service
          uri: lb://order-service
          predicates:
            - Path=/api/order/**
          filters:
            - StripPrefix=1
            - name: RequestRateLimiter
              args:
                redis-rate-limiter.replenishRate: 500
                redis-rate-limiter.burstCapacity: 1000
            - name: CircuitBreaker
              args:
                name: orderCB
                fallbackUri: forward:/fallback/order
      loadbalancer:
        retry:
          enabled: true
          maxRetriesOnSameServiceInstance: 1
          maxRetriesOnNextServiceInstance: 3
      globalcors:
        cors-configurations:
          '[/**]':
            allowedOrigins: "*"
            allowedMethods: "*"

resilience4j:
  circuitbreaker:
    instances:
      orderCB:
        slidingWindowSize: 20
        failureRateThreshold: 40
        waitDurationInOpenState: 30s
  timelimiter:
    instances:
      orderCB:
        timeoutDuration: 5s

3.2 动态路由配置源

@Component
public class NacosDynamicRouteConfig {

    private final ConfigService configService;
    private final ApplicationEventPublisher publisher;
    private final GatewayProperties gatewayProperties;

    public NacosDynamicRouteConfig(
            ConfigService configService,
            ApplicationEventPublisher publisher,
            GatewayProperties gatewayProperties) {
        this.configService = configService;
        this.publisher = publisher;
        this.gatewayProperties = gatewayProperties;
    }

    @PostConstruct
    public void init() throws Exception {
        String dataId = "gateway-routes.json";
        configService.addListener(dataId, "DEFAULT_GROUP", new Listener() {
            @Override
            public Executor getExecutor() {
                return Executors.newSingleThreadExecutor();
            }

            @Override
            public void receiveConfigInfo(String configInfo) {
                List<RouteDefinition> routes = JSON.parseArray(
                    configInfo, RouteDefinition.class);
                gatewayProperties.setRoutes(routes);
                publisher.publishEvent(new RefreshRoutesEvent(this));
            }
        });
    }
}

四、高级实践

4.1 自适应负载均衡器

public class AdaptiveLoadBalancer implements ReactorServiceInstanceLoadBalancer {

    private final String serviceId;
    private final ObjectProvider<ServiceInstanceListSupplier> supplierProvider;
    private final MeterRegistry meterRegistry;

    public AdaptiveLoadBalancer(
            String serviceId,
            ObjectProvider<ServiceInstanceListSupplier> supplierProvider,
            MeterRegistry meterRegistry) {
        this.serviceId = serviceId;
        this.supplierProvider = supplierProvider;
        this.meterRegistry = meterRegistry;
    }

    @Override
    public Mono<Response<ServiceInstance>> choose(Request request) {
        ServiceInstanceListSupplier supplier = supplierProvider.getIfAvailable();
        return supplier.get(request).next().map(instances -> {
            List<ServiceInstance> scored = instances.stream()
                    .map(instance -> {
                        double score = calculateScore(instance);
                        return new AbstractMap.SimpleEntry<>(instance, score);
                    })
                    .sorted(Map.Entry.comparingByValue())
                    .collect(Collectors.toList());

            if (scored.isEmpty()) {
                return Response.empty();
            }

            double totalScore = scored.stream()
                    .mapToDouble(Map.Entry::getValue).sum();
            double random = ThreadLocalRandom.current().nextDouble(totalScore);
            double cumulative = 0.0;

            for (Map.Entry<ServiceInstance, Double> entry : scored) {
                cumulative += entry.getValue();
                if (random <= cumulative) {
                    ServiceInstance chosen = entry.getKey();
                    recordMetrics(chosen);
                    return Response.just(chosen);
                }
            }

            return Response.just(scored.get(0).getKey());
        });
    }

    private double calculateScore(ServiceInstance instance) {
        Map<String, String> metadata = instance.getMetadata();
        double baseScore = 100.0;

        String cpuStr = metadata.get("cpuUsage");
        if (cpuStr != null) {
            double cpu = Double.parseDouble(cpuStr);
            baseScore -= cpu * 50;
        }

        String rtStr = metadata.get("avgResponseTime");
        if (rtStr != null) {
            double rt = Double.parseDouble(rtStr);
            if (rt > 1000) {
                baseScore -= 30;
            } else if (rt > 500) {
                baseScore -= 15;
            }
        }

        String activeStr = metadata.get("activeConnections");
        if (activeStr != null) {
            int active = Integer.parseInt(activeStr);
            baseScore -= Math.min(active * 2, 40);
        }

        return Math.max(baseScore, 1.0);
    }

    private void recordMetrics(ServiceInstance instance) {
        meterRegistry.counter(
            "gateway.loadbalancer.choose",
            "serviceId", serviceId,
            "instance", instance.getHost() + ":" + instance.getPort()
        ).increment();
    }
}

4.2 熔断与过载联动

@Component
public class OverloadProtectionManager {

    private final Map<String, AtomicInteger> overloadCounters = new ConcurrentHashMap<>();
    private final Map<String, Boolean> circuitStates = new ConcurrentHashMap<>();
    private static final int OVERLOAD_THRESHOLD = 10;
    private static final long RECOVERY_WINDOW_MS = 30000;

    public boolean tryAcquire(String serviceId) {
        if (isCircuitOpen(serviceId)) {
            if (tryRecovery(serviceId)) {
                return true;
            }
            return false;
        }

        return true;
    }

    public void recordOverload(String serviceId) {
        AtomicInteger counter = overloadCounters
                .computeIfAbsent(serviceId, k -> new AtomicInteger(0));
        int count = counter.incrementAndGet();
        if (count >= OVERLOAD_THRESHOLD) {
            openCircuit(serviceId);
        }
    }

    public void recordSuccess(String serviceId) {
        AtomicInteger counter = overloadCounters
                .computeIfAbsent(serviceId, k -> new AtomicInteger(0));
        counter.set(Math.max(0, counter.get() - 1));
    }

    private boolean isCircuitOpen(String serviceId) {
        return circuitStates.getOrDefault(serviceId, false);
    }

    private void openCircuit(String serviceId) {
        circuitStates.put(serviceId, true);
        ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor();
        executor.schedule(() -> {
            circuitStates.put(serviceId, false);
            overloadCounters.getOrDefault(serviceId, new AtomicInteger(0)).set(0);
        }, RECOVERY_WINDOW_MS, TimeUnit.MILLISECONDS);
        executor.shutdown();
    }

    private boolean tryRecovery(String serviceId) {
        return false;
    }
}

4.3 路由刷新风暴防护

@Component
public class RouteRefreshThrottle {

    private final AtomicLong lastRefreshTime = new AtomicLong(0);
    private final Queue<Long> refreshHistory = new ConcurrentLinkedQueue<>();
    private static final long MIN_INTERVAL_MS = 2000;
    private static final int MAX_REFRESHES_PER_MINUTE = 10;

    public boolean allowRefresh() {
        long now = System.currentTimeMillis();
        long last = lastRefreshTime.get();

        if (now - last < MIN_INTERVAL_MS) {
            log.warn("路由刷新频率过高,拒绝本次刷新,距离上次仅 {}ms", now - last);
            return false;
        }

        refreshHistory.add(now);
        long oneMinuteAgo = now - 60000;
        while (!refreshHistory.isEmpty()
                && refreshHistory.peek() < oneMinuteAgo) {
            refreshHistory.poll();
        }

        if (refreshHistory.size() >= MAX_REFRESHES_PER_MINUTE) {
            log.warn("路由刷新超过每分钟上限,当前 {} 次/分钟", refreshHistory.size());
            return false;
        }

        lastRefreshTime.set(now);
        return true;
    }
}

4.4 优雅降级Filter

@Component
public class GracefulDegradationFilter implements GlobalFilter, Ordered {

    private final OverloadProtectionManager protectionManager;

    public GracefulDegradationFilter(
            OverloadProtectionManager protectionManager) {
        this.protectionManager = protectionManager;
    }

    @Override
    public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) {
        String serviceId = extractServiceId(exchange);
        if (serviceId == null) {
            return chain.filter(exchange);
        }

        if (!protectionManager.tryAcquire(serviceId)) {
            return degradeResponse(exchange, serviceId, "circuit_open");
        }

        return chain.filter(exchange)
                .doOnSuccess(unused ->
                    protectionManager.recordSuccess(serviceId))
                .onErrorResume(throwable -> {
                    if (isOverloadError(throwable)) {
                        protectionManager.recordOverload(serviceId);
                        return degradeResponse(exchange, serviceId, "overload");
                    }
                    return Mono.error(throwable);
                });
    }

    private Mono<Void> degradeResponse(
            ServerWebExchange exchange, String serviceId, String reason) {
        exchange.getResponse().setStatusCode(HttpStatus.SERVICE_UNAVAILABLE);
        exchange.getResponse().getHeaders()
                .setContentType(MediaType.APPLICATION_JSON);

        byte[] body = String.format(
            "{\"code\":503,\"message\":\"服务[%s]暂不可用,原因: %s\",\"timestamp\":%d}",
            serviceId, reason, System.currentTimeMillis()
        ).getBytes(StandardCharsets.UTF_8);

        return exchange.getResponse()
                .writeWith(Mono.just(exchange.getResponse()
                        .bufferFactory().wrap(body)));
    }

    private String extractServiceId(ServerWebExchange exchange) {
        Route route = exchange.getAttribute(
            ServerWebExchangeUtils.GATEWAY_ROUTE_ATTR);
        if (route == null) {
            return null;
        }
        String uri = route.getUri().toString();
        if (uri.startsWith("lb://")) {
            return uri.substring(5);
        }
        return null;
    }

    private boolean isOverloadError(Throwable t) {
        if (t instanceof TimeoutException) {
            return true;
        }
        if (t instanceof CircuitBreakerException) {
            return true;
        }
        String message = t.getMessage();
        return message != null && (
            message.contains("timeout")
            || message.contains("circuit")
            || message.contains("overload")
        );
    }

    @Override
    public int getOrder() {
        return Ordered.LOWEST_PRECEDENCE - 100;
    }
}

五、最佳实践

实践要点 说明 推荐度
自适应负载均衡 根据CPU/RT/活跃连接数动态评分,避免流量倾斜到慢节点 ⭐⭐⭐⭐⭐
路由变更节流 限制每分钟路由刷新次数,防止Nacos推送风暴 ⭐⭐⭐⭐⭐
熔断过载联动 熔断器打开时联动过载保护管理器,全局控制流量准入 ⭐⭐⭐⭐
优雅降级响应 降级时返回标准化JSON,包含服务名和原因,便于排查 ⭐⭐⭐⭐
指标埋点 负载均衡选择、限流命中、熔断打开等关键事件埋点到Prometheus ⭐⭐⭐⭐
渐进恢复 熔断恢复时从小流量逐步放开,避免瞬间涌入压垮服务 ⭐⭐⭐⭐⭐

六、总结

高并发场景下,Spring Cloud Gateway的动态路由、负载均衡与过载防护三者必须联动配合才能发挥最大效能。动态路由需要节流保护防止配置风暴,负载均衡需要自适应评分实现智能调度,过载防护需要全局协调实现优雅降级。

本文通过自适应负载均衡器、路由刷新节流器、过载保护管理器、优雅降级Filter等组件,构建了一套完整的流量治理体系。这套方案的关键在于"感知-决策-执行"的闭环——实时感知实例状态和流量变化,动态调整路由和负载均衡决策,在超过承受能力时优雅执行降级动作。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐