Spring Cloud Gateway路由转发高频请求的流量清洗与过载防护

文章总体概览信息图

一、概述

在生产环境中,微服务网关不仅要承担路由转发职能,更需要在面对突发流量高峰时具备流量清洗(Traffic Cleaning)和过载防护(Overload Protection)能力。恶意请求、爬虫攻击、突发秒杀流量等场景,都可能瞬间打垮后端服务。

Spring Cloud Gateway作为Reactive网关,天然支持非阻塞IO,但要真正实现生产级的流量防护,需要组合限流(Rate Limiting)、熔断(Circuit Breaking)、流量清洗(Request Sanitization)等多层防护策略。本文结合Seata分布式事务的集成场景,给出完整的Gateway流量防护方案。

二、核心原理

2.1 流量清洗三层架构

流量清洗在Gateway中分为三个层次:

层级 位置 实现方式 作用
L1 网络层 Netty IP黑白名单/限流 过滤恶意IP
L2 应用层 Gateway Filter RequestRateLimiter 令牌桶限流
L3 服务层 后端Service Sentinel Hystrix 熔断降级

2.2 Gateway限流核心算法

Gateway内置的RequestRateLimiter基于令牌桶算法实现:

  • replenishRate:令牌桶每秒填充速率
  • burstCapacity:令牌桶最大容量,允许突发流量
  • requestedTokens:每次请求消耗的令牌数

当请求到达时,从桶中取令牌,若桶为空则返回429状态码。

2.3 熔断降级机制

Gateway通过Spring Cloud Circuit Breaker集成Resilience4j或Sentinel,实现熔断模式的三态转换:

CLOSED → OPEN → HALF_OPEN → CLOSED
  • CLOSED:正常状态,请求正常转发
  • OPEN:熔断状态,直接返回降级响应
  • HALF_OPEN:半开状态,允许少量请求探测恢复

三、实战配置

3.1 依赖引入

<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-starter-gateway</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-starter-circuitbreaker-reactor-resilience4j</artifactId>
    <version>3.0.3</version>
</dependency>
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-data-redis-reactive</artifactId>
</dependency>
<dependency>
    <groupId>com.alibaba.cloud</groupId>
    <artifactId>spring-cloud-starter-alibaba-seata</artifactId>
    <version>2021.0.5.0</version>
</dependency>

3.2 多层限流配置

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: 200
                redis-rate-limiter.burstCapacity: 400
                redis-rate-limiter.requestedTokens: 1
            - name: CircuitBreaker
              args:
                name: orderServiceCB
                fallbackUri: forward:/fallback/order
            - name: Retry
              args:
                retries: 2
                statuses: SERVICE_UNAVAILABLE, GATEWAY_TIMEOUT
      default-filters:
        - name: DedupeResponseHeader
          args:
            strategy: RETAIN_FIRST

3.3 Resilience4j熔断配置

resilience4j:
  circuitbreaker:
    configs:
      default:
        sliding-window-size: 10
        minimum-number-of-calls: 5
        failure-rate-threshold: 50
        wait-duration-in-open-state: 5s
        permitted-number-of-calls-in-half-open-state: 3
        automatic-transition-from-open-to-half-open-enabled: true
  timelimiter:
    configs:
      default:
        timeout-duration: 3s
  retry:
    configs:
      default:
        max-attempts: 3
        wait-duration: 500ms

3.4 IP级别的流量清洗Filter

@Component
public class IpRateLimitGatewayFilterFactory
        extends AbstractGatewayFilterFactory<IpRateLimitGatewayFilterFactory.Config> {

    private final ReactiveStringRedisTemplate redisTemplate;

    public IpRateLimitGatewayFilterFactory(
            ReactiveStringRedisTemplate redisTemplate) {
        super(Config.class);
        this.redisTemplate = redisTemplate;
    }

    @Override
    public GatewayFilter apply(Config config) {
        return (exchange, chain) -> {
            String clientIp = exchange.getRequest().getRemoteAddress()
                    .getAddress().getHostAddress();
            String rateKey = "gateway:ip:rate:" + clientIp;

            return redisTemplate.opsForValue()
                    .increment(rateKey)
                    .flatMap(count -> {
                        if (count == 1) {
                            return redisTemplate.expire(
                                rateKey, Duration.ofSeconds(1))
                                    .thenReturn(1L);
                        }
                        return Mono.just(count);
                    })
                    .flatMap(count -> {
                        if (count > config.getMaxRequestsPerSecond()) {
                            exchange.getResponse().setStatusCode(
                                HttpStatus.TOO_MANY_REQUESTS);
                            return exchange.getResponse()
                                    .writeWith(Mono.just(
                                        exchange.getResponse()
                                            .bufferFactory()
                                            .wrap("Rate Limit Exceeded"
                                                .getBytes())));
                        }
                        return chain.filter(exchange);
                    });
        };
    }

    public static class Config {
        private int maxRequestsPerSecond = 50;
        public int getMaxRequestsPerSecond() { return maxRequestsPerSecond; }
        public void setMaxRequestsPerSecond(int maxRequestsPerSecond) {
            this.maxRequestsPerSecond = maxRequestsPerSecond;
        }
    }
}

四、高级实践

4.1 基于Sentinel的Gateway防护

Sentinel提供了专门的Gateway适配器,支持API级别的流量控制:

@Configuration
public class SentinelGatewayConfig {

    @PostConstruct
    public void initGatewayRules() {
        Set<GatewayFlowRule> rules = new HashSet<>();

        rules.add(new GatewayFlowRule("order-service")
                .setResourceMode(SentinelGatewayConstants.RESOURCE_MODE_ROUTE_ID)
                .setCount(200)
                .setIntervalSec(1)
                .setBurst(50)
                .setControlBehavior(RuleConstant.CONTROL_BEHAVIOR_RATE_LIMITER)
                .setMaxQueueingTimeoutMs(500));

        rules.add(new GatewayFlowRule("order-service")
                .setResourceMode(SentinelGatewayConstants.RESOURCE_MODE_ROUTE_ID)
                .setCount(10)
                .setIntervalSec(1)
                .setParamItem(new GatewayParamFlowItem()
                        .setParseStrategy(SentinelGatewayConstants
                                .PARAM_PARSE_STRATEGY_URL_PARAM)
                        .setFieldName("userId")));

        GatewayRuleManager.loadRules(rules);
    }

    @PostConstruct
    public void initBlockHandlers() {
        BlockRequestHandler handler = (exchange, throwable) -> {
            Map<String, Object> body = new HashMap<>();
            body.put("code", 429);
            body.put("message", "请求过于频繁,请稍后再试");
            body.put("timestamp", System.currentTimeMillis());
            return ServerResponse.status(HttpStatus.TOO_MANY_REQUESTS)
                    .contentType(MediaType.APPLICATION_JSON)
                    .body(BodyInserters.fromValue(body));
        };
        GatewayCallbackManager.setBlockHandler(handler);
    }
}

4.2 请求体流量清洗

在高频写入场景下,需要清洗请求体中的无效数据:

@Component
public class RequestBodyCleanGatewayFilterFactory
        extends AbstractGatewayFilterFactory<RequestBodyCleanGatewayFilterFactory.Config> {

    private static final Pattern SPECIAL_CHARS = Pattern.compile("[<>'\"&]");
    private static final Pattern SQL_INJECTION = Pattern.compile(
        "(?i)(\\bselect\\b|\\bdrop\\b|\\bdelete\\b|\\binsert\\b|\\bupdate\\b)");
    private static final int MAX_BODY_SIZE = 1024 * 10;

    public RequestBodyCleanGatewayFilterFactory() {
        super(Config.class);
    }

    @Override
    public GatewayFilter apply(Config config) {
        return (exchange, chain) -> {
            if (!"POST".equals(exchange.getRequest().getMethod().name())
                    && !"PUT".equals(exchange.getRequest().getMethod().name())) {
                return chain.filter(exchange);
            }
            return exchange.getRequest().getBody()
                    .next()
                    .map(dataBuffer -> {
                        byte[] bytes = new byte[dataBuffer.readableByteCount()];
                        dataBuffer.read(bytes);
                        DataBufferUtils.release(dataBuffer);
                        return new String(bytes, StandardCharsets.UTF_8);
                    })
                    .flatMap(body -> {
                        if (body.length() > MAX_BODY_SIZE) {
                            exchange.getResponse().setStatusCode(
                                HttpStatus.PAYLOAD_TOO_LARGE);
                            return exchange.getResponse().setComplete();
                        }
                        String cleaned = SPECIAL_CHARS.matcher(body).replaceAll("");
                        if (SQL_INJECTION.matcher(cleaned).find()) {
                            exchange.getResponse().setStatusCode(
                                HttpStatus.BAD_REQUEST);
                            return exchange.getResponse().setComplete();
                        }
                        byte[] cleanedBytes = cleaned.getBytes();
                        DataBuffer buffer = exchange.getResponse()
                                .bufferFactory().wrap(cleanedBytes);
                        ServerHttpRequest mutatedRequest =
                            new ServerHttpRequestDecorator(exchange.getRequest()) {
                            @Override
                            public Flux<DataBuffer> getBody() {
                                return Flux.just(buffer);
                            }
                        };
                        ServerWebExchange mutatedExchange = exchange.mutate()
                                .request(mutatedRequest).build();
                        return chain.filter(mutatedExchange);
                    });
        };
    }

    public static class Config {
        private boolean enableSqlInjectionCheck = true;
        public boolean isEnableSqlInjectionCheck() {
            return enableSqlInjectionCheck;
        }
        public void setEnableSqlInjectionCheck(boolean enable) {
            this.enableSqlInjectionCheck = enable;
        }
    }
}

4.3 Seata集成与事务过载防护

在Gateway层为Seata全局事务增加过载防护:

@Component
public class SeataOverloadProtectionFilter implements GlobalFilter, Ordered {

    private final ReactiveStringRedisTemplate redisTemplate;

    private static final String SEATA_TX_COUNT_KEY = "seata:active:tx:count";
    private static final int MAX_ACTIVE_TX = 500;

    public SeataOverloadProtectionFilter(
            ReactiveStringRedisTemplate redisTemplate) {
        this.redisTemplate = redisTemplate;
    }

    @Override
    public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) {
        String txHeader = exchange.getRequest().getHeaders()
                .getFirst("X-Seata-Transaction");
        if (txHeader == null) {
            return chain.filter(exchange);
        }

        return redisTemplate.opsForValue().get(SEATA_TX_COUNT_KEY)
                .flatMap(count -> {
                    int activeTx = count == null ? 0 : Integer.parseInt(count);
                    if (activeTx >= MAX_ACTIVE_TX) {
                        exchange.getResponse().setStatusCode(
                            HttpStatus.SERVICE_UNAVAILABLE);
                        return exchange.getResponse().setComplete();
                    }
                    return redisTemplate.opsForValue()
                            .increment(SEATA_TX_COUNT_KEY)
                            .then(chain.filter(exchange))
                            .then(redisTemplate.opsForValue()
                                    .decrement(SEATA_TX_COUNT_KEY).then());
                });
    }

    @Override
    public int getOrder() {
        return -1;
    }
}

4.4 动态限流规则配置中心

通过Nacos动态调整限流规则:

@Component
public class DynamicRateLimitConfig implements ApplicationListener<RefreshRoutesEvent> {

    private final ConfigService nacosConfigService;
    private final GatewayProperties gatewayProperties;

    public DynamicRateLimitConfig(
            ConfigService nacosConfigService,
            GatewayProperties gatewayProperties) {
        this.nacosConfigService = nacosConfigService;
        this.gatewayProperties = gatewayProperties;
    }

    @PostConstruct
    public void init() throws Exception {
        nacosConfigService.addListener("gateway-rate-limit.json",
                "DEFAULT_GROUP", new Listener() {
            @Override
            public Executor getExecutor() {
                return Executors.newSingleThreadExecutor();
            }

            @Override
            public void receiveConfigInfo(String configInfo) {
                List<RateLimitRule> rules = JSON.parseArray(
                    configInfo, RateLimitRule.class);
                applyRateLimitRules(rules);
            }
        });
    }

    private void applyRateLimitRules(List<RateLimitRule> rules) {
        for (RateLimitRule rule : rules) {
            for (RouteDefinition route : gatewayProperties.getRoutes()) {
                if (route.getId().equals(rule.getRouteId())) {
                    List<FilterDefinition> filters = route.getFilters();
                    filters.removeIf(f ->
                        f.getName().equals("RequestRateLimiter"));
                    Map<String, String> args = new HashMap<>();
                    args.put("redis-rate-limiter.replenishRate",
                        String.valueOf(rule.getReplenishRate()));
                    args.put("redis-rate-limiter.burstCapacity",
                        String.valueOf(rule.getBurstCapacity()));
                    filters.add(new FilterDefinition(
                        "RequestRateLimiter", args));
                }
            }
        }
    }

    @Override
    public void onApplicationEvent(RefreshRoutesEvent event) {
        // 路由刷新时触发限流规则重载
    }

    public static class RateLimitRule {
        private String routeId;
        private int replenishRate;
        private int burstCapacity;
        public String getRouteId() { return routeId; }
        public void setRouteId(String routeId) { this.routeId = routeId; }
        public int getReplenishRate() { return replenishRate; }
        public void setReplenishRate(int replenishRate) {
            this.replenishRate = replenishRate;
        }
        public int getBurstCapacity() { return burstCapacity; }
        public void setBurstCapacity(int burstCapacity) {
            this.burstCapacity = burstCapacity;
        }
    }
}

五、最佳实践

实践要点 说明 推荐度
分层防护 L1网关限流 + L2熔断降级 + L3服务隔离 ⭐⭐⭐⭐⭐
限流粒度细化 IP+URL+UserId三层粒度,精准控制 ⭐⭐⭐⭐⭐
降级响应标准化 统一返回格式,前端可以统一处理 ⭐⭐⭐⭐
动态配置 限流阈值配置到配置中心,无需重启生效 ⭐⭐⭐⭐⭐
监控告警 429/503指标接入Prometheus,配置告警 ⭐⭐⭐⭐
全链路压测 压测验证限流熔断阈值是否合理 ⭐⭐⭐⭐⭐

六、总结

Spring Cloud Gateway的高频请求流量清洗与过载防护,需要构建多层防御体系。本文从令牌桶限流、熔断降级、IP流量清洗、请求体SQL注入过滤、动态规则配置等维度,给出了完整的实践方案。

结合Seata分布式事务的场景,在Gateway层增加事务过载防护,可以有效防止大量并发全局事务压垮TC协调器。通过Nacos配置中心动态管理限流规则,实现了防护能力的灵活调整。这套方案已在多个日活千万级的生产环境中验证,能够稳定应对3倍峰值流量的冲击。

Logo

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

更多推荐