Java熔断降级案例如何开发

wen java案例 29

本文目录导读:

Java熔断降级案例如何开发

  1. 使用 Resilience4j 实现熔断降级
  2. 使用 Sentinel 实现熔断降级
  3. 使用 Hystrix 实现熔断降级(已进入维护模式)
  4. 完整示例:分布式服务熔断降级
  5. 最佳实践建议
  6. 测试用例

我来为你提供一个完整的Java熔断降级案例开发指南,包含多种实现方式。

使用 Resilience4j 实现熔断降级

1 Maven依赖

<dependency>
    <groupId>io.github.resilience4j</groupId>
    <artifactId>resilience4j-spring-boot2</artifactId>
    <version>2.1.0</version>
</dependency>
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-aop</artifactId>
</dependency>

2 配置文件

resilience4j:
  circuitbreaker:
    configs:
      default:
        sliding-window-size: 10
        minimum-number-of-calls: 5
        failure-rate-threshold: 50
        wait-duration-in-open-state: 10s
        permitted-number-of-calls-in-half-open-state: 3
        record-exceptions:
          - java.lang.Exception
    instances:
      userService:
        base-config: default

3 实现代码

@Service
public class UserService {
    private static final Logger logger = LoggerFactory.getLogger(UserService.class);
    @CircuitBreaker(name = "userService", fallbackMethod = "getUserFallback")
    public User getUser(Long userId) {
        // 模拟调用远程服务
        return callRemoteService(userId);
    }
    public User getUserFallback(Long userId, Throwable t) {
        logger.error("获取用户信息失败,使用降级数据", t);
        return new User(userId, "默认用户", "未知");
    }
    private User callRemoteService(Long userId) {
        // 模拟远程调用,随机失败
        if (Math.random() > 0.5) {
            throw new RuntimeException("远程服务调用失败");
        }
        return new User(userId, "张三", "zhangsan@example.com");
    }
}

使用 Sentinel 实现熔断降级

1 Maven依赖

<dependency>
    <groupId>com.alibaba.csp</groupId>
    <artifactId>sentinel-core</artifactId>
    <version>1.8.6</version>
</dependency>
<dependency>
    <groupId>com.alibaba.csp</groupId>
    <artifactId>sentinel-annotation-aspectj</artifactId>
    <version>1.8.6</version>
</dependency>

2 实现代码

@Component
public class OrderService {
    private static final Logger logger = LoggerFactory.getLogger(OrderService.class);
    @SentinelResource(
        value = "createOrder",
        fallback = "createOrderFallback",
        blockHandler = "createOrderBlockHandler"
    )
    public Order createOrder(OrderRequest request) {
        // 模拟创建订单
        if (request.getAmount() > 10000) {
            throw new RuntimeException("订单金额超出限制");
        }
        return doCreateOrder(request);
    }
    // 业务异常降级
    public Order createOrderFallback(OrderRequest request, Throwable t) {
        logger.error("创建订单失败,使用降级策略", t);
        return Order.buildFallbackOrder(request);
    }
    // 限流降级
    public Order createOrderBlockHandler(OrderRequest request, BlockException ex) {
        logger.warn("请求被限流,使用降级策略");
        return Order.buildBlockedOrder(request);
    }
}
// 配置类
@Configuration
public class SentinelConfig {
    @PostConstruct
    public void initRules() {
        // 熔断降级规则
        DegradeRule degradeRule = new DegradeRule("createOrder")
            .setGrade(RuleConstant.DEGRADE_GRADE_RT)  // RT降级
            .setCount(100)  // 平均响应时间 > 100ms
            .setTimeWindow(10)  // 熔断时长10秒
            .setMinRequestAmount(5);  // 最小请求数
        DegradeRuleManager.loadRules(Collections.singletonList(degradeRule));
        // 限流规则
        FlowRule flowRule = new FlowRule("createOrder")
            .setGrade(RuleConstant.FLOW_GRADE_QPS)  // QPS限流
            .setCount(10)  // 最大QPS为10
            .setControlBehavior(RuleConstant.CONTROL_BEHAVIOR_DEFAULT);
        FlowRuleManager.loadRules(Collections.singletonList(flowRule));
    }
}

使用 Hystrix 实现熔断降级(已进入维护模式)

1 Maven依赖

<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-starter-netflix-hystrix</artifactId>
    <version>2.2.10.RELEASE</version>
</dependency>

2 实现代码

@Service
public class PaymentService {
    private static final Logger logger = LoggerFactory.getLogger(PaymentService.class);
    @HystrixCommand(
        fallbackMethod = "processPaymentFallback",
        commandProperties = {
            @HystrixProperty(name = "execution.isolation.thread.timeoutInMilliseconds", value = "1000"),
            @HystrixProperty(name = "circuitBreaker.requestVolumeThreshold", value = "10"),
            @HystrixProperty(name = "circuitBreaker.errorThresholdPercentage", value = "50"),
            @HystrixProperty(name = "circuitBreaker.sleepWindowInMilliseconds", value = "5000")
        }
    )
    public PaymentResult processPayment(PaymentRequest request) {
        // 模拟支付处理
        return callPaymentGateway(request);
    }
    public PaymentResult processPaymentFallback(PaymentRequest request) {
        logger.warn("支付服务降级,请求ID: {}", request.getOrderId());
        return PaymentResult.fallbackResult(request);
    }
    private PaymentResult callPaymentGateway(PaymentRequest request) {
        // 模拟支付网关调用
        if (Math.random() > 0.7) {
            throw new HystrixRuntimeException(
                HystrixRuntimeException.FailureType.COMMAND_EXCEPTION,
                null, 
                "支付网关超时"
            );
        }
        return new PaymentResult(true, "支付成功");
    }
}

完整示例:分布式服务熔断降级

1 业务场景

// 订单处理服务
@Service
public class OrderProcessingService {
    private final UserService userService;
    private final InventoryService inventoryService;
    private final PaymentService paymentService;
    @Autowired
    public OrderProcessingService(
            UserService userService,
            InventoryService inventoryService,
            PaymentService paymentService) {
        this.userService = userService;
        this.inventoryService = inventoryService;
        this.paymentService = paymentService;
    }
    @CircuitBreaker(name = "orderProcess", fallbackMethod = "processOrderFallback")
    @Bulkhead(name = "orderProcess", type = Bulkhead.Type.SEMAPHORE)
    @Retry(name = "orderProcess", fallbackMethod = "retryFallback")
    public OrderResult processOrder(Order order) {
        // 1. 验证用户
        User user = userService.getUser(order.getUserId());
        if (user == null) {
            return OrderResult.failure("用户不存在");
        }
        // 2. 检查库存
        InventoryCheckResult inventory = inventoryService.checkInventory(
            order.getProductId(), order.getQuantity());
        if (!inventory.isAvailable()) {
            return OrderResult.failure("库存不足");
        }
        // 3. 处理支付
        PaymentResult payment = paymentService.processPayment(
            new PaymentRequest(order.getOrderId(), order.getAmount()));
        if (!payment.isSuccess()) {
            // 释放库存
            inventoryService.releaseInventory(order.getProductId(), order.getQuantity());
            return OrderResult.failure("支付失败");
        }
        return OrderResult.success("订单处理成功");
    }
    public OrderResult processOrderFallback(Order order, Throwable t) {
        logger.error("订单处理降级,订单号: {}, 原因: {}", 
            order.getOrderId(), t.getMessage());
        return OrderResult.fallback("系统繁忙,请稍后重试");
    }
    public OrderResult retryFallback(Order order, Throwable t) {
        logger.error("重试失败,订单号: {}", order.getOrderId());
        return OrderResult.failure("服务暂时不可用");
    }
}

2 熔断器监控

@Component
public class CircuitBreakerMonitor {
    private final CircuitBreakerRegistry circuitBreakerRegistry;
    @Autowired
    public CircuitBreakerMonitor(CircuitBreakerRegistry circuitBreakerRegistry) {
        this.circuitBreakerRegistry = circuitBreakerRegistry;
        startMonitoring();
    }
    public void startMonitoring() {
        ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1);
        scheduler.scheduleAtFixedRate(() -> {
            circuitBreakerRegistry.getAllCircuitBreakers().forEach(cb -> {
                CircuitBreaker.Metrics metrics = cb.getMetrics();
                logger.info("Circuit Breaker: {}, State: {}, " +
                    "Failure Rate: {}, Success Count: {}, Failure Count: {}",
                    cb.getName(),
                    cb.getState(),
                    metrics.getFailureRate(),
                    metrics.getNumberOfSuccessfulCalls(),
                    metrics.getNumberOfFailedCalls()
                );
            });
        }, 0, 10, TimeUnit.SECONDS);
    }
}

最佳实践建议

1 配置优化

resilience4j:
  circuitbreaker:
    configs:
      default:
        # 滑动窗口大小
        sliding-window-size: 20
        # 最小请求数
        minimum-number-of-calls: 10
        # 失败率阈值
        failure-rate-threshold: 50
        # 慢调用阈值
        slow-call-rate-threshold: 50
        slow-call-duration-threshold: 5s
        # 熔断持续时间
        wait-duration-in-open-state: 30s
        # 半开状态尝试次数
        permitted-number-of-calls-in-half-open-state: 3
        # 异常记录
        record-exceptions:
          - java.lang.Exception
        ignore-exceptions:
          - java.lang.IllegalArgumentException

2 降级策略设计

public interface FallbackStrategy<T> {
    T handleFallback(Object... args);
}
@Component
public class CacheFallbackStrategy implements FallbackStrategy<User> {
    private final RedisTemplate<String, User> redisTemplate;
    @Override
    public User handleFallback(Object... args) {
        Long userId = (Long) args[0];
        // 从缓存获取降级数据
        User cachedUser = redisTemplate.opsForValue().get("user:" + userId);
        if (cachedUser != null) {
            return cachedUser;
        }
        // 返回默认空对象
        return new User(userId, null, null);
    }
}

3 健康检查接口

@RestController
@RequestMapping("/actuator")
public class CircuitBreakerController {
    @GetMapping("/circuit-breakers")
    public Map<String, Object> getCircuitBreakersStatus() {
        CircuitBreakerRegistry registry = CircuitBreakerRegistry.ofDefaults();
        Map<String, Object> status = new HashMap<>();
        registry.getAllCircuitBreakers().forEach(cb -> {
            Map<String, Object> cbStatus = new HashMap<>();
            cbStatus.put("state", cb.getState());
            cbStatus.put("metrics", cb.getMetrics());
            status.put(cb.getName(), cbStatus);
        });
        return status;
    }
}

测试用例

@SpringBootTest
public class CircuitBreakerTest {
    @Autowired
    private OrderProcessingService orderService;
    @Test
    public void testCircuitBreaker() {
        // 模拟大量失败请求触发熔断
        for (int i = 0; i < 20; i++) {
            Order order = new Order();
            order.setUserId((long) (i % 5));
            order.setAmount(1000);
            OrderResult result = orderService.processOrder(order);
            if (i > 10) {
                // 熔断后的降级响应
                assertTrue(result.isFallback());
            }
        }
    }
    @Test
    public void testCircuitBreakerStateTransition() {
        CircuitBreaker circuitBreaker = circuitBreakerRegistry
            .circuitBreaker("orderProcess");
        // 初始状态应为 CLOSED
        assertEquals(CircuitBreaker.State.CLOSED, circuitBreaker.getState());
        // 触发熔断条件
        for (int i = 0; i < 20; i++) {
            try {
                orderService.processOrder(createInvalidOrder());
            } catch (Exception ignored) {}
        }
        // 检查是否进入 OPEN 状态
        assertEquals(CircuitBreaker.State.OPEN, circuitBreaker.getState());
    }
}

这个完整的案例涵盖了主流的熔断降级实现方式,包括配置、代码实现、监控和测试,建议根据实际项目需求选择合适的框架。

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