本文目录导读:

下面我为您提供一个完整的Java熔断降级案例,使用Spring Cloud Circuit Breaker + Resilience4j实现。
项目依赖配置
<dependencies>
<!-- Spring Boot Starter -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<!-- Spring Cloud Circuit Breaker -->
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-starter-circuitbreaker-resilience4j</artifactId>
</dependency>
<!-- Resilience4j 配置支持 -->
<dependency>
<groupId>io.github.resilience4j</groupId>
<artifactId>resilience4j-spring-boot2</artifactId>
</dependency>
<!-- Actuator 监控 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
</dependency>
</dependencies>
application.yml 配置
server:
port: 8080
spring:
application:
name: circuit-breaker-demo
# Resilience4j 配置
resilience4j:
circuitbreaker:
instances:
userService: # 针对具体服务的配置
registerHealthIndicator: true
slidingWindowSize: 10 # 滑动窗口大小
minimumNumberOfCalls: 5 # 最少调用次数,达到后才开始计算
permittedNumberOfCallsInHalfOpenState: 3 # 半开状态允许的调用次数
automaticTransitionFromOpenToHalfOpenEnabled: true
waitDurationInOpenState: 5s # 打开状态持续时间
failureRateThreshold: 60 # 失败率阈值(百分比)
eventConsumerBufferSize: 10
recordExceptions:
- java.lang.Exception
ignoreExceptions:
- org.springframework.web.client.RestClientException
timelimiter:
instances:
userService:
timeoutDuration: 3s # 超时时间
cancelRunningFuture: true
bulkhead:
instances:
userService:
maxConcurrentCalls: 10 # 最大并发调用数
maxWaitDuration: 10ms # 最大等待时间
management:
endpoints:
web:
exposure:
include: "*"
health:
circuitbreakers:
enabled: true
服务层实现
import io.github.resilience4j.circuitbreaker.annotation.CircuitBreaker;
import io.github.resilience4j.bulkhead.annotation.Bulkhead;
import io.github.resilience4j.ratelimiter.annotation.RateLimiter;
import io.github.resilience4j.timelimiter.annotation.TimeLimiter;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionStage;
@Service
public class UserService {
private static final Logger logger = LoggerFactory.getLogger(UserService.class);
private final ExternalUserApiClient externalClient;
public UserService(ExternalUserApiClient externalClient) {
this.externalClient = externalClient;
}
/**
* 使用熔断器获取用户信息
*/
@CircuitBreaker(name = "userService", fallbackMethod = "getUserFallback")
@Bulkhead(name = "userService", fallbackMethod = "getUserBulkheadFallback")
public User getUser(String userId) {
logger.info("尝试获取用户: {}", userId);
return externalClient.fetchUser(userId);
}
/**
* 熔断器的降级方法
*/
public User getUserFallback(String userId, Exception exception) {
logger.error("熔断器触发,降级处理 userId: {}, 原因: {}", userId, exception.getMessage());
// 构建缓存中的用户或者默认用户
return getCachedUser(userId);
}
/**
* 隔离舱降级方法
*/
public User getUserBulkheadFallback(String userId, Exception exception) {
logger.error("隔离舱触发,降级处理 userId: {}, 原因: {}", userId, exception.getMessage());
return getDefaultUser(userId);
}
/**
* 使用限流器
*/
@RateLimiter(name = "userService", fallbackMethod = "rateLimiterFallback")
public User getUserWithRateLimit(String userId) {
return externalClient.fetchUser(userId);
}
public User rateLimiterFallback(String userId, Exception exception) {
logger.error("限流触发,降级处理 userId: {}, 原因: {}", userId, exception.getMessage());
return getDefaultUser(userId);
}
/**
* 异步调用用户服务(带超时限制)
*/
@CircuitBreaker(name = "userService", fallbackMethod = "asyncGetUserFallback")
@TimeLimiter(name = "userService")
public CompletionStage<User> asyncGetUser(String userId) {
return CompletableFuture.supplyAsync(() -> externalClient.fetchUser(userId));
}
public CompletionStage<User> asyncGetUserFallback(String userId, Exception exception) {
logger.error("异步调用熔断,userId: {}, 原因: {}", userId, exception.getMessage());
return CompletableFuture.completedFuture(getDefaultUser(userId));
}
// 降级辅助方法
private User getCachedUser(String userId) {
// 从本地缓存获取用户信息
User cachedUser = LocalCache.getUser(userId);
if (cachedUser != null) {
return cachedUser;
}
return getDefaultUser(userId);
}
private User getDefaultUser(String userId) {
return User.builder()
.id(userId)
.name("默认用户")
.email("default@example.com")
.status("DEFAULT")
.build();
}
}
Controller 层实现
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.*;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionStage;
@RestController
@RequestMapping("/api/users")
public class UserController {
@Autowired
private UserService userService;
/**
* 同步获取用户信息(带熔断)
*/
@GetMapping("/{userId}")
public ResponseEntity<User> getUser(@PathVariable String userId) {
try {
User user = userService.getUser(userId);
return ResponseEntity.ok(user);
} catch (Exception e) {
// 即使熔断器失败,也返回降级数据
User fallbackUser = userService.getUserFallback(userId, e);
return ResponseEntity.status(200).body(fallbackUser);
}
}
/**
* 获取用户信息(带限流)
*/
@GetMapping("/limited/{userId}")
public ResponseEntity<User> getUserWithRateLimit(@PathVariable String userId) {
User user = userService.getUserWithRateLimit(userId);
return ResponseEntity.ok(user);
}
/**
* 异步获取用户信息
*/
@GetMapping("/async/{userId}")
public CompletableFuture<ResponseEntity<User>> asyncGetUser(@PathVariable String userId) {
return userService.asyncGetUser(userId)
.thenApply(ResponseEntity::ok)
.toCompletableFuture();
}
/**
* 测试熔断器状态
*/
@GetMapping("/circuit/status")
public ResponseEntity<CircuitBreakerStatus> getCircuitBreakerStatus() {
return ResponseEntity.ok(CircuitBreakerStatusService.getStatus());
}
}
模拟外部API客户端
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
import org.springframework.web.client.RestTemplate;
import java.util.Random;
@Component
public class ExternalUserApiClient {
private static final Logger logger = LoggerFactory.getLogger(ExternalUserApiClient.class);
private final RestTemplate restTemplate;
private final Random random = new Random();
public ExternalUserApiClient() {
this.restTemplate = new RestTemplate();
}
public User fetchUser(String userId) {
logger.info("调用外部API获取用户: {}", userId);
// 模拟外部API调用,有 60% 的概率失败
if (random.nextInt(100) < 60) {
throw new RuntimeException("外部服务调用失败");
}
// 模拟网络延迟
try {
Thread.sleep(100);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
// 模拟正常的服务响应
return User.builder()
.id(userId)
.name("用户" + userId)
.email("user" + userId + "@example.com")
.status("ACTIVE")
.build();
}
}
熔断器状态监控工具类
import io.github.resilience4j.circuitbreaker.CircuitBreaker;
import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import java.util.HashMap;
import java.util.Map;
@Component
public class CircuitBreakerStatusService {
@Autowired
private CircuitBreakerRegistry circuitBreakerRegistry;
public Map<String, Object> getStatus() {
Map<String, Object> statusMap = new HashMap<>();
CircuitBreaker circuitBreaker = circuitBreakerRegistry.circuitBreaker("userService");
CircuitBreaker.Metrics metrics = circuitBreaker.getMetrics();
statusMap.put("state", circuitBreaker.getState().name());
statusMap.put("failureRate", metrics.getFailureRate());
statusMap.put("successRate", metrics.getSuccessRate());
statusMap.put("calls", metrics.getNumberOfSuccessfulCalls());
statusMap.put("failedCalls", metrics.getNumberOfFailedCalls());
statusMap.put("bufferedCalls", metrics.getNumberOfBufferedCalls());
statusMap.put("notPermittedCalls", metrics.getNumberOfNotPermittedCalls());
return statusMap;
}
}
用户实体类
import lombok.Builder;
import lombok.Data;
@Data
@Builder
public class User {
private String id;
private String name;
private String email;
private String status;
}
测试示例
@RestController
@RequestMapping("/test")
public class CircuitBreakerTestController {
@Autowired
private UserService userService;
/**
* 测试熔断器行为
* 连续调用此接口,观察熔断器状态变化
*/
@GetMapping("/breaker-test/{userId}")
public ResponseEntity<Map<String, Object>> testCircuitBreaker(@PathVariable String userId) {
Map<String, Object> result = new HashMap<>();
long startTime = System.currentTimeMillis();
try {
User user = userService.getUser(userId);
result.put("success", true);
result.put("user", user);
} catch (Exception e) {
result.put("success", false);
result.put("error", e.getMessage());
}
long endTime = System.currentTimeMillis();
result.put("duration", (endTime - startTime) + "ms");
return ResponseEntity.ok(result);
}
}
配置类(可选)
import io.github.resilience4j.circuitbreaker.CircuitBreaker;
import io.github.resilience4j.circuitbreaker.CircuitBreakerConfig;
import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.time.Duration;
@Configuration
public class Resilience4jConfig {
@Bean
public CircuitBreakerRegistry circuitBreakerRegistry() {
CircuitBreakerConfig config = CircuitBreakerConfig.custom()
.failureRateThreshold(60) // 失败率阈值
.waitDurationInOpenState(Duration.ofSeconds(5)) // 打开状态持续时间
.slidingWindowSize(10) // 滑动窗口大小
.minimumNumberOfCalls(5) // 最少调用次数
.permittedNumberOfCallsInHalfOpenState(3) // 半开状态允许的调用
.build();
return CircuitBreakerRegistry.of(config);
}
}
使用说明
- 运行项目:启动Spring Boot应用
- 测试熔断器:多次调用
/api/users/{userId}接口,当失败率达到60%时,熔断器会打开 - 查看监控:访问
/actuator/health查看熔断器健康状态 - 测试限流:快速多次调用
/api/users/limited/{userId}观察限流效果
这个案例展示了:
- 熔断器(CircuitBreaker):失败率超过阈值时自动熔断
- 降级(Fallback):熔断后返回默认数据
- 限流(RateLimiter):控制请求速率
- 隔离舱(Bulkhead):限制并发调用数
- 超时控制(TimeLimiter):限制调用时间