本文目录导读:

我来提供一个完整的Java实现API网关案例,包含核心功能:路由转发、认证、限流、日志等。
项目结构
api-gateway/
├── pom.xml
├── src/main/java/com/gateway/
│ ├── ApiGatewayApplication.java
│ ├── config/
│ │ ├── GatewayConfig.java
│ │ └── RedisConfig.java
│ ├── filter/
│ │ ├── AuthenticationFilter.java
│ │ ├── RateLimitFilter.java
│ │ ├── LogFilter.java
│ │ └── RouteFilter.java
│ ├── model/
│ │ ├── RouteDefinition.java
│ │ └── ApiResponse.java
│ ├── service/
│ │ ├── RouteService.java
│ │ ├── AuthenticationService.java
│ │ └── RateLimitService.java
│ └── controller/
│ └── GatewayController.java
pom.xml
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.gateway</groupId>
<artifactId>api-gateway</artifactId>
<version>1.0.0</version>
<packaging>jar</packaging>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>2.7.0</version>
<relativePath/>
</parent>
<properties>
<java.version>11</java.version>
<spring-cloud.version>2021.0.3</spring-cloud.version>
</properties>
<dependencies>
<!-- Spring Boot Web -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<!-- Spring Cloud Gateway -->
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-starter-gateway</artifactId>
</dependency>
<!-- Redis for rate limiting -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis-reactive</artifactId>
</dependency>
<!-- JWT for authentication -->
<dependency>
<groupId>io.jsonwebtoken</groupId>
<artifactId>jjwt-api</artifactId>
<version>0.11.5</version>
</dependency>
<dependency>
<groupId>io.jsonwebtoken</groupId>
<artifactId>jjwt-impl</artifactId>
<version>0.11.5</version>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>io.jsonwebtoken</groupId>
<artifactId>jjwt-jackson</artifactId>
<version>0.11.5</version>
<scope>runtime</scope>
</dependency>
<!-- Lombok -->
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
<!-- Test -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
<dependencyManagement>
<dependencies>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-dependencies</artifactId>
<version>${spring-cloud.version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>
</project>
主应用类
package com.gateway;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.cloud.gateway.route.RouteLocator;
import org.springframework.cloud.gateway.route.builder.RouteLocatorBuilder;
import org.springframework.context.annotation.Bean;
@SpringBootApplication
public class ApiGatewayApplication {
public static void main(String[] args) {
SpringApplication.run(ApiGatewayApplication.class, args);
}
@Bean
public RouteLocator customRouteLocator(RouteLocatorBuilder builder) {
return builder.routes()
// 用户服务路由
.route("user-service", r -> r
.path("/api/user/**")
.filters(f -> f
.stripPrefix(2)
.addRequestHeader("X-Gateway", "api-gateway"))
.uri("http://localhost:8081/"))
// 订单服务路由
.route("order-service", r -> r
.path("/api/order/**")
.filters(f -> f
.stripPrefix(2)
.addRequestHeader("X-Gateway", "api-gateway"))
.uri("http://localhost:8082/"))
// 商品服务路由
.route("product-service", r -> r
.path("/api/product/**")
.filters(f -> f
.stripPrefix(2)
.addRequestHeader("X-Gateway", "api-gateway"))
.uri("http://localhost:8083/"))
.build();
}
}
路由配置
package com.gateway.config;
import org.springframework.cloud.gateway.filter.ratelimit.KeyResolver;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import reactor.core.publisher.Mono;
@Configuration
public class GatewayConfig {
/**
* 基于IP的限流键解析器
*/
@Bean
public KeyResolver ipKeyResolver() {
return exchange -> {
String ip = exchange.getRequest().getRemoteAddress().getAddress().getHostAddress();
return Mono.just(ip);
};
}
/**
* 基于用户的限流键解析器(需要从JWT中获取用户信息)
*/
@Bean
public KeyResolver userKeyResolver() {
return exchange -> {
String userId = exchange.getRequest().getHeaders()
.getFirst("X-User-Id");
if (userId == null) {
userId = "anonymous";
}
return Mono.just(userId);
};
}
}
过滤器实现
1 认证过滤器
package com.gateway.filter;
import io.jsonwebtoken.Claims;
import io.jsonwebtoken.Jwts;
import io.jsonwebtoken.security.Keys;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.cloud.gateway.filter.GatewayFilter;
import org.springframework.cloud.gateway.filter.GatewayFilterChain;
import org.springframework.core.Ordered;
import org.springframework.http.HttpStatus;
import org.springframework.http.server.reactive.ServerHttpRequest;
import org.springframework.http.server.reactive.ServerHttpResponse;
import org.springframework.stereotype.Component;
import org.springframework.web.server.ServerWebExchange;
import reactor.core.publisher.Mono;
import javax.crypto.SecretKey;
import java.nio.charset.StandardCharsets;
@Component
public class AuthenticationFilter implements GatewayFilter, Ordered {
@Value("${jwt.secret}")
private String jwtSecret;
@Value("${jwt.expiration}")
private long jwtExpiration;
@Override
public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) {
ServerHttpRequest request = exchange.getRequest();
String path = request.getPath().toString();
// 白名单路径(不需要认证)
if (isWhitelisted(path)) {
return chain.filter(exchange);
}
// 获取token
String token = extractToken(request);
if (token == null) {
return unauthorized(exchange, "缺少认证令牌");
}
try {
// 验证token
SecretKey key = Keys.hmacShaKeyFor(jwtSecret.getBytes(StandardCharsets.UTF_8));
Claims claims = Jwts.parserBuilder()
.setSigningKey(key)
.build()
.parseClaimsJws(token)
.getBody();
// 将用户信息添加到请求头
ServerHttpRequest mutatedRequest = request.mutate()
.header("X-User-Id", claims.getSubject())
.header("X-User-Role", String.valueOf(claims.get("role")))
.build();
return chain.filter(exchange.mutate().request(mutatedRequest).build());
} catch (Exception e) {
return unauthorized(exchange, "无效的令牌");
}
}
@Override
public int getOrder() {
return -100; // 最高优先级
}
private boolean isWhitelisted(String path) {
// 白名单路径
String[] whitelist = {
"/api/auth/login",
"/api/auth/register",
"/api/public"
};
for (String item : whitelist) {
if (path.startsWith(item)) {
return true;
}
}
return false;
}
private String extractToken(ServerHttpRequest request) {
String bearer = request.getHeaders().getFirst("Authorization");
if (bearer != null && bearer.startsWith("Bearer ")) {
return bearer.substring(7);
}
return null;
}
private Mono<Void> unauthorized(ServerWebExchange exchange, String message) {
ServerHttpResponse response = exchange.getResponse();
response.setStatusCode(HttpStatus.UNAUTHORIZED);
response.getHeaders().add("Content-Type", "application/json");
String body = String.format("{\"code\":401,\"message\":\"%s\"}", message);
return response.writeWith(Mono.just(exchange.getResponse()
.bufferFactory().wrap(body.getBytes())));
}
}
2 限流过滤器
package com.gateway.filter;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.gateway.filter.GatewayFilter;
import org.springframework.cloud.gateway.filter.GatewayFilterChain;
import org.springframework.core.Ordered;
import org.springframework.data.redis.core.ReactiveStringRedisTemplate;
import org.springframework.http.HttpStatus;
import org.springframework.http.server.reactive.ServerHttpResponse;
import org.springframework.stereotype.Component;
import org.springframework.web.server.ServerWebExchange;
import reactor.core.publisher.Mono;
import java.time.Duration;
import java.time.LocalDateTime;
@Component
public class RateLimitFilter implements GatewayFilter, Ordered {
@Autowired
private ReactiveStringRedisTemplate redisTemplate;
// 限制配置
private static final int MAX_REQUESTS_PER_MINUTE = 100;
private static final String RATE_LIMIT_PREFIX = "rate_limit:";
@Override
public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) {
String clientIp = getClientIp(exchange);
String key = RATE_LIMIT_PREFIX + clientIp;
String currentMinute = LocalDateTime.now().format(java.time.format.DateTimeFormatter.ofPattern("yyyyMMddHHmm"));
String redisKey = key + ":" + currentMinute;
return redisTemplate.opsForValue()
.increment(redisKey)
.flatMap(count -> {
// 设置过期时间(60秒)
if (count == 1) {
return redisTemplate.expire(redisKey, Duration.ofSeconds(60))
.thenReturn(count);
}
return Mono.just(count);
})
.flatMap(count -> {
if (count > MAX_REQUESTS_PER_MINUTE) {
return rateLimitExceeded(exchange);
}
return chain.filter(exchange);
});
}
@Override
public int getOrder() {
return -90; // 在认证过滤器之后
}
private String getClientIp(ServerWebExchange exchange) {
String xForwardedFor = exchange.getRequest()
.getHeaders().getFirst("X-Forwarded-For");
if (xForwardedFor != null && !xForwardedFor.isEmpty()) {
return xForwardedFor.split(",")[0].trim();
}
return exchange.getRequest().getRemoteAddress().getAddress().getHostAddress();
}
private Mono<Void> rateLimitExceeded(ServerWebExchange exchange) {
ServerHttpResponse response = exchange.getResponse();
response.setStatusCode(HttpStatus.TOO_MANY_REQUESTS);
response.getHeaders().add("Content-Type", "application/json");
String body = "{\"code\":429,\"message\":\"请求过于频繁,请稍后再试\"}";
return response.writeWith(Mono.just(response.bufferFactory().wrap(body.getBytes())));
}
}
3 日志过滤器
package com.gateway.filter;
import lombok.extern.slf4j.Slf4j;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.cloud.gateway.filter.GatewayFilter;
import org.springframework.cloud.gateway.filter.GatewayFilterChain;
import org.springframework.core.Ordered;
import org.springframework.http.server.reactive.ServerHttpRequest;
import org.springframework.stereotype.Component;
import org.springframework.util.StopWatch;
import org.springframework.web.server.ServerWebExchange;
import reactor.core.publisher.Mono;
@Slf4j
@Component
public class LogFilter implements GatewayFilter, Ordered {
private static final Logger logger = LoggerFactory.getLogger(LogFilter.class);
@Override
public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) {
StopWatch stopWatch = new StopWatch();
stopWatch.start();
ServerHttpRequest request = exchange.getRequest();
String requestId = java.util.UUID.randomUUID().toString();
// 添加请求ID
exchange.getRequest().mutate()
.header("X-Request-ID", requestId)
.build();
// 记录请求日志
logger.info("请求开始 [{}] {} {}", requestId,
request.getMethodValue(), request.getURI());
return chain.filter(exchange)
.doFinally(signalType -> {
stopWatch.stop();
logger.info("请求完成 [{}] 状态: {}, 耗时: {}ms",
requestId,
exchange.getResponse().getStatusCode(),
stopWatch.getTotalTimeMillis());
});
}
@Override
public int getOrder() {
return -50;
}
}
服务实现
1 动态路由服务
package com.gateway.service;
import com.gateway.model.RouteDefinition;
import lombok.Data;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.gateway.route.RouteDefinitionWriter;
import org.springframework.stereotype.Service;
import org.springframework.web.reactive.function.server.ServerResponse;
import reactor.core.publisher.Mono;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@Service
public class RouteService {
@Autowired
private RouteDefinitionWriter routeDefinitionWriter;
@Autowired
private RouteDefinitionLocator routeDefinitionLocator;
private final Map<String, String> routeCache = new ConcurrentHashMap<>();
public Mono<Boolean> addRoute(RouteDefinition routeDef) {
return routeDefinitionWriter
.save(Mono.just(convertToRouteDefinition(routeDef)))
.thenReturn(true);
}
public Mono<Boolean> deleteRoute(String routeId) {
return routeDefinitionWriter
.delete(Mono.just(routeId))
.thenReturn(true);
}
public Mono<Boolean> updateRoute(RouteDefinition routeDef) {
return routeDefinitionWriter
.save(Mono.just(convertToRouteDefinition(routeDef)))
.thenReturn(true);
}
private org.springframework.cloud.gateway.route.RouteDefinition
convertToRouteDefinition(RouteDefinition routeDef) {
org.springframework.cloud.gateway.route.RouteDefinition def =
new org.springframework.cloud.gateway.route.RouteDefinition();
def.setId(routeDef.getId());
org.springframework.cloud.gateway.support.UriUtils uri =
org.springframework.cloud.gateway.support.UriUtils.convertToUri(routeDef.getUrl());
def.setUri(uri);
// 设置路径谓词
org.springframework.cloud.gateway.handler.predicate.PathRoutePredicateFactory.PathConfig pathConfig =
new org.springframework.cloud.gateway.handler.predicate.PathRoutePredicateFactory.PathConfig();
pathConfig.setPatterns(routeDef.getPathPatterns());
org.springframework.cloud.gateway.filter.FilterDefinition filterDef =
new org.springframework.cloud.gateway.filter.FilterDefinition();
filterDef.setName("StripPrefix");
filterDef.addArg("parts", String.valueOf(routeDef.getStripPrefix()));
def.getPredicates().add(
new org.springframework.cloud.gateway.handler.predicate.PredicateDefinition(
"Path=" + String.join(",", routeDef.getPathPatterns())
)
);
def.getFilters().add(filterDef);
return def;
}
}
2 认证服务
package com.gateway.service;
import io.jsonwebtoken.Jwts;
import io.jsonwebtoken.SignatureAlgorithm;
import io.jsonwebtoken.security.Keys;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import javax.crypto.SecretKey;
import java.nio.charset.StandardCharsets;
import java.util.Date;
import java.util.HashMap;
import java.util.Map;
@Service
public class AuthenticationService {
@Value("${jwt.secret}")
private String jwtSecret;
@Value("${jwt.expiration}")
private long jwtExpiration;
@Value("${jwt.issuer}")
private String jwtIssuer;
public String generateToken(String userId, String role) {
SecretKey key = Keys.hmacShaKeyFor(jwtSecret.getBytes(StandardCharsets.UTF_8));
Map<String, Object> claims = new HashMap<>();
claims.put("role", role);
Date now = new Date();
Date expiryDate = new Date(now.getTime() + jwtExpiration);
return Jwts.builder()
.setClaims(claims)
.setSubject(userId)
.setIssuer(jwtIssuer)
.setIssuedAt(now)
.setExpiration(expiryDate)
.signWith(key, SignatureAlgorithm.HS256)
.compact();
}
public boolean validateToken(String token) {
try {
SecretKey key = Keys.hmacShaKeyFor(jwtSecret.getBytes(StandardCharsets.UTF_8));
Jwts.parserBuilder()
.setSigningKey(key)
.build()
.parseClaimsJws(token);
return true;
} catch (Exception e) {
return false;
}
}
public String getUserIdFromToken(String token) {
SecretKey key = Keys.hmacShaKeyFor(jwtSecret.getBytes(StandardCharsets.UTF_8));
return Jwts.parserBuilder()
.setSigningKey(key)
.build()
.parseClaimsJws(token)
.getBody()
.getSubject();
}
}
3 限流服务
package com.gateway.service;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.ReactiveStringRedisTemplate;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Mono;
import java.time.Duration;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
@Service
public class RateLimitService {
@Autowired
private ReactiveStringRedisTemplate redisTemplate;
private static final String RATE_LIMIT_PREFIX = "gateway:ratelimit:";
private static final DateTimeFormatter FORMATTER =
DateTimeFormatter.ofPattern("yyyyMMddHHmm");
public Mono<Boolean> checkRateLimit(String key, long maxRequests, Duration window) {
String redisKey = buildRedisKey(key);
return redisTemplate.opsForValue()
.increment(redisKey)
.flatMap(count -> {
if (count == 1) {
return redisTemplate.expire(redisKey, window)
.thenReturn(true);
}
return Mono.just(true);
})
.map(count -> count <= maxRequests);
}
private String buildRedisKey(String key) {
String minute = LocalDateTime.now().format(FORMATTER);
return RATE_LIMIT_PREFIX + key + ":" + minute;
}
public Mono<Void> resetRateLimit(String key) {
String redisKey = buildRedisKey(key);
return redisTemplate.delete(redisKey).then();
}
}
控制器
package com.gateway.controller;
import com.gateway.model.ApiResponse;
import com.gateway.model.RouteDefinition;
import com.gateway.service.AuthenticationService;
import com.gateway.service.RouteService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;
import reactor.core.publisher.Mono;
import java.util.HashMap;
import java.util.Map;
@RestController
@RequestMapping("/api/gateway")
public class GatewayController {
@Autowired
private AuthenticationService authenticationService;
@Autowired
private RouteService routeService;
@PostMapping("/auth/token")
public Mono<ApiResponse> generateToken(@RequestBody LoginRequest loginRequest) {
// 这里应该验证用户信息,这里简化处理
String token = authenticationService.generateToken(
loginRequest.getUserId(),
loginRequest.getRole()
);
Map<String, String> data = new HashMap<>();
data.put("token", token);
return Mono.just(ApiResponse.success(data));
}
@PostMapping("/routes")
public Mono<ApiResponse> addRoute(@RequestBody RouteDefinition routeDefinition) {
return routeService.addRoute(routeDefinition)
.map(success -> ApiResponse.success(null, "路由添加成功"));
}
@DeleteMapping("/routes/{routeId}")
public Mono<ApiResponse> deleteRoute(@PathVariable String routeId) {
return routeService.deleteRoute(routeId)
.map(success -> ApiResponse.success(null, "路由删除成功"));
}
@PutMapping("/routes")
public Mono<ApiResponse> updateRoute(@RequestBody RouteDefinition routeDefinition) {
return routeService.updateRoute(routeDefinition)
.map(success -> ApiResponse.success(null, "路由更新成功"));
}
// 内部类
static class LoginRequest {
private String userId;
private String role;
public String getUserId() { return userId; }
public void setUserId(String userId) { this.userId = userId; }
public String getRole() { return role; }
public void setRole(String role) { this.role = role; }
}
}
模型类
package com.gateway.model;
import lombok.Data;
import java.util.List;
@Data
public class RouteDefinition {
private String id;
private String name;
private String url; // 目标服务地址
private List<String> pathPatterns; // 路径匹配模式
private int stripPrefix = 1; // 去除前缀数量
private boolean authRequired = true; // 是否需要认证
private Long rateLimit; // 限流阈值
}
package com.gateway.model;
import lombok.Data;
@Data
public class ApiResponse<T> {
private int code;
private String message;
private T data;
private long timestamp;
public static <T> ApiResponse<T> success(T data) {
ApiResponse<T> response = new ApiResponse<>();
response.setCode(200);
response.setMessage("success");
response.setData(data);
response.setTimestamp(System.currentTimeMillis());
return response;
}
public static <T> ApiResponse<T> success(T data, String message) {
ApiResponse<T> response = new ApiResponse<>();
response.setCode(200);
response.setMessage(message);
response.setData(data);
response.setTimestamp(System.currentTimeMillis());
return response;
}
public static <T> ApiResponse<T> error(int code, String message) {
ApiResponse<T> response = new ApiResponse<>();
response.setCode(code);
response.setMessage(message);
response.setData(null);
response.setTimestamp(System.currentTimeMillis());
return response;
}
}
配置文件
# application.yml
server:
port: 8080
spring:
application:
name: api-gateway
redis:
host: localhost
port: 6379
database: 0
timeout: 5000ms
lettuce:
pool:
max-active: 8
max-idle: 8
min-idle: 0
cloud:
gateway:
# 全局过滤器配置
default-filters:
- DedupeResponseHeader=Access-Control-Allow-Origin
globalcors:
cors-configurations:
'[/**]':
allowed-origins: "*"
allowed-methods: "*"
allowed-headers: "*"
allow-credentials: true
# HTTP客户端配置
codec:
max-in-memory-size: 10MB
# JWT配置
jwt:
secret: your-secret-key-with-at-least-32-characters
expiration: 86400000 # 24小时
issuer: api-gateway
# 日志配置
logging:
level:
root: INFO
org.springframework.cloud.gateway: DEBUG
com.gateway: DEBUG
使用示例
启动网关
# 确保Redis已启动 redis-server # 启动API网关 mvn spring-boot:run
测试API
# 1. 获取Token
curl -X POST http://localhost:8080/api/gateway/auth/token \
-H "Content-Type: application/json" \
-d '{"userId": "user123", "role": "ADMIN"}'
# 2. 调用受保护的API
curl -X GET http://localhost:8080/api/user/info \
-H "Authorization: Bearer eyJhbGciOiJIUzI1NiJ9..."
# 3. 动态添加路由
curl -X POST http://localhost:8080/api/gateway/routes \
-H "Content-Type: application/json" \
-d '{
"id": "new-service",
"name": "新服务",
"url": "http://localhost:8090",
"pathPatterns": ["/api/new/**"],
"stripPrefix": 2,
"authRequired": true
}'
# 4. 测试限流
for i in {1..110}; do
curl -X GET http://localhost:8080/api/user/info \
-H "Authorization: Bearer your-token"
done
扩展功能建议
- 负载均衡 - 集成Netflix Ribbon或Spring Cloud LoadBalancer
- 服务发现 - 集成Nacos、Consul或Zookeeper
- 熔断降级 - 集成Hystrix或Resilience4j
- 灰度发布 - 根据权重或header路由到不同版本服务
- 监控告警 - 集成Prometheus + Grafana
- API文档 - 集成Swagger聚合
- 安全增强 - 集成OAuth2.0、RBAC权限控制
- 消息推送 - 集成WebSocket和SSE支持
这个案例提供了一个完整的API网关实现,涵盖了路由转发、认证、限流、日志等核心功能,可以快速集成到生产环境中。