Java熔断降级案例

wen java案例 1

本文目录导读:

Java熔断降级案例

  1. 项目依赖配置
  2. application.yml 配置
  3. 服务层实现
  4. Controller 层实现
  5. 模拟外部API客户端
  6. 熔断器状态监控工具类
  7. 用户实体类
  8. 测试示例
  9. 配置类(可选)
  10. 使用说明

下面我为您提供一个完整的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);
    }
}

使用说明

  1. 运行项目:启动Spring Boot应用
  2. 测试熔断器:多次调用 /api/users/{userId} 接口,当失败率达到60%时,熔断器会打开
  3. 查看监控:访问 /actuator/health 查看熔断器健康状态
  4. 测试限流:快速多次调用 /api/users/limited/{userId} 观察限流效果

这个案例展示了:

  • 熔断器(CircuitBreaker):失败率超过阈值时自动熔断
  • 降级(Fallback):熔断后返回默认数据
  • 限流(RateLimiter):控制请求速率
  • 隔离舱(Bulkhead):限制并发调用数
  • 超时控制(TimeLimiter):限制调用时间

抱歉,评论功能暂时关闭!