Spring Cloud Gateway路由转发高频请求的流量清洗与过载防护
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倍峰值流量的冲击。
更多推荐




所有评论(0)