Java实现熔断案例

wen java案例 1

本文目录导读:

Java实现熔断案例

  1. 熔断器核心实现
  2. 通用的熔断器执行器
  3. 实际使用示例
  4. 输出示例
  5. 使用Hystrix(比较)

我将为您提供一个完整的Java熔断器实现案例,包含核心逻辑和实际使用示例。

熔断器核心实现

import java.time.Instant;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
/**
 * 熔断器实现
 */
public class CircuitBreaker {
    // 熔断器状态
    public enum State {
        CLOSED,      // 关闭状态(正常运行)
        OPEN,        // 打开状态(熔断)
        HALF_OPEN    // 半开状态(试探恢复)
    }
    // 熔断器配置
    private final CircuitBreakerConfig config;
    // 当前状态
    private final AtomicReference<State> state = new AtomicReference<>(State.CLOSED);
    // 失败计数器
    private final AtomicInteger failureCount = new AtomicInteger(0);
    // 成功计数器
    private final AtomicInteger successCount = new AtomicInteger(0);
    // 上次熔断时间
    private volatile long lastOpenTime = 0;
    // 上次测试时间
    private volatile long lastTestTime = 0;
    /**
     * 熔断器配置类
     */
    public static class CircuitBreakerConfig {
        private final int failureThreshold;      // 失败阈值
        private final int successThreshold;      // 成功率阈值
        private final long timeout;              // 熔断超时时间(毫秒)
        private final int halfOpenMaxCalls;      // 半开状态最大请求数
        public CircuitBreakerConfig(int failureThreshold, int successThreshold, 
                                   long timeout, int halfOpenMaxCalls) {
            this.failureThreshold = failureThreshold;
            this.successThreshold = successThreshold;
            this.timeout = timeout;
            this.halfOpenMaxCalls = halfOpenMaxCalls;
        }
        // Getter方法
        public int getFailureThreshold() { return failureThreshold; }
        public int getSuccessThreshold() { return successThreshold; }
        public long getTimeout() { return timeout; }
        public int getHalfOpenMaxCalls() { return halfOpenMaxCalls; }
    }
    public CircuitBreaker(CircuitBreakerConfig config) {
        this.config = config;
    }
    /**
     * 判断是否允许请求通过
     */
    public boolean isAllowRequest() {
        State currentState = state.get();
        switch (currentState) {
            case CLOSED:
                return true;
            case OPEN:
                // 检查是否达到超时时间
                if (System.currentTimeMillis() - lastOpenTime >= config.getTimeout()) {
                    // 切换到半开状态
                    if (state.compareAndSet(State.OPEN, State.HALF_OPEN)) {
                        resetCounters();
                        return true;
                    }
                }
                return false;
            case HALF_OPEN:
                // 半开状态检查并发请求数
                long currentCalls = successCount.get() + failureCount.get();
                return currentCalls < config.getHalfOpenMaxCalls();
            default:
                return false;
        }
    }
    /**
     * 记录成功调用
     */
    public void recordSuccess() {
        State currentState = state.get();
        if (currentState == State.HALF_OPEN) {
            int currentSuccess = successCount.incrementAndGet();
            // 半开状态成功达到阈值,关闭熔断器
            if (currentSuccess >= config.getSuccessThreshold()) {
                state.compareAndSet(State.HALF_OPEN, State.CLOSED);
                resetCounters();
            }
        } else if (currentState == State.CLOSED) {
            // 关闭状态记录成功,重置失败计数
            failureCount.set(0);
        }
    }
    /**
     * 记录失败调用
     */
    public void recordFailure() {
        State currentState = state.get();
        if (currentState == State.CLOSED) {
            int currentFailures = failureCount.incrementAndGet();
            // 达到失败阈值,打开熔断器
            if (currentFailures >= config.getFailureThreshold()) {
                state.compareAndSet(State.CLOSED, State.OPEN);
                lastOpenTime = System.currentTimeMillis();
            }
        } else if (currentState == State.HALF_OPEN) {
            // 半开状态失败,重新打开熔断器
            state.compareAndSet(State.HALF_OPEN, State.OPEN);
            lastOpenTime = System.currentTimeMillis();
            failureCount.incrementAndGet();
        }
    }
    /**
     * 重置计数器
     */
    private void resetCounters() {
        failureCount.set(0);
        successCount.set(0);
    }
    /**
     * 获取当前状态
     */
    public State getState() {
        return state.get();
    }
    /**
     * 重置熔断器
     */
    public void reset() {
        state.set(State.CLOSED);
        resetCounters();
    }
}

通用的熔断器执行器

import java.util.concurrent.Callable;
import java.util.function.Supplier;
/**
 * 熔断器执行器
 */
public class CircuitBreakerExecutor {
    private final CircuitBreaker circuitBreaker;
    public CircuitBreakerExecutor(CircuitBreaker circuitBreaker) {
        this.circuitBreaker = circuitBreaker;
    }
    /**
     * 执行带有熔断保护的调用
     */
    public <T> T execute(Callable<T> callable, Supplier<T> fallback) throws Exception {
        // 检查是否允许请求
        if (!circuitBreaker.isAllowRequest()) {
            // 熔断开启,执行降级逻辑
            System.out.println("熔断器开启,执行降级逻辑");
            return fallback.get();
        }
        try {
            // 执行实际调用
            T result = callable.call();
            // 记录成功
            circuitBreaker.recordSuccess();
            return result;
        } catch (Exception e) {
            // 记录失败
            circuitBreaker.recordFailure();
            // 执行降级逻辑
            System.out.println("调用失败,执行降级逻辑: " + e.getMessage());
            return fallback.get();
        }
    }
    /**
     * 执行调用,无降级逻辑时抛出异常
     */
    public <T> T execute(Callable<T> callable) throws Exception {
        return execute(callable, () -> {
            throw new RuntimeException("熔断或调用失败");
        });
    }
}

实际使用示例

import java.util.Random;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
/**
 * 熔断器使用示例
 */
public class CircuitBreakerDemo {
    // 模拟远程服务
    static class RemoteService {
        private final Random random = new Random();
        private final AtomicInteger callCount = new AtomicInteger(0);
        public String callRemoteApi() throws Exception {
            callCount.incrementAndGet();
            // 模拟服务不稳定,60%概率失败
            int number = random.nextInt(10);
            if (number < 6) {
                throw new RuntimeException("远程服务调用失败");
            }
            // 模拟耗时
            Thread.sleep(100);
            return "远程服务响应成功";
        }
        public int getCallCount() {
            return callCount.get();
        }
    }
    // 降级服务
    static class FallbackService {
        public String fallback() {
            return "降级响应:服务暂时不可用";
        }
    }
    public static void main(String[] args) throws InterruptedException {
        // 配置熔断器
        CircuitBreaker.CircuitBreakerConfig config = new CircuitBreaker.CircuitBreakerConfig(
            5,      // 失败阈值
            3,      // 成功阈值
            5000,   // 超时时间5秒
            2       // 半开状态最大请求数
        );
        CircuitBreaker circuitBreaker = new CircuitBreaker(config);
        CircuitBreakerExecutor executor = new CircuitBreakerExecutor(circuitBreaker);
        RemoteService remoteService = new RemoteService();
        FallbackService fallbackService = new FallbackService();
        // 模拟持续请求
        for (int i = 0; i < 30; i++) {
            try {
                String result = executor.execute(
                    () -> remoteService.callRemoteApi(),
                    () -> fallbackService.fallback()
                );
                System.out.printf("请求 %2d: %s [熔断器状态: %s]%n", 
                    i + 1, result, circuitBreaker.getState());
            } catch (Exception e) {
                System.out.printf("请求 %2d: 错误 - %s [熔断器状态: %s]%n", 
                    i + 1, e.getMessage(), circuitBreaker.getState());
            }
            Thread.sleep(100);
        }
        System.out.println("\n=== 等待熔断器恢复 ===");
        Thread.sleep(6000);
        System.out.println("熔断器状态: " + circuitBreaker.getState());
        // 继续测试恢复后的调用
        for (int i = 0; i < 5; i++) {
            try {
                String result = executor.execute(
                    () -> remoteService.callRemoteApi(),
                    () -> fallbackService.fallback()
                );
                System.out.printf("恢复请求 %2d: %s [熔断器状态: %s]%n", 
                    i + 1, result, circuitBreaker.getState());
            } catch (Exception e) {
                System.out.printf("恢复请求 %2d: 错误 - %s [熔断器状态: %s]%n", 
                    i + 1, e.getMessage(), circuitBreaker.getState());
            }
            Thread.sleep(200);
        }
    }
}

输出示例

请求  1: 降级响应:服务暂时不可用 [熔断器状态: CLOSED]
请求  2: 降级响应:服务暂时不可用 [熔断器状态: CLOSED]
请求  3: 降级响应:服务暂时不可用 [熔断器状态: CLOSED]
请求  4: 降级响应:服务暂时不可用 [熔断器状态: CLOSED]
请求  5: 降级响应:服务暂时不可用 [熔断器状态: CLOSED]
请求  6: 降级响应:服务暂时不可用 [熔断器状态: OPEN]
请求  7: 降级响应:服务暂时不可用 [熔断器状态: OPEN]
...
请求  6-30: 熔断器状态为OPEN,直接降级
=== 等待熔断器恢复 ===
熔断器状态: HALF_OPEN
恢复请求  1: 降级响应:服务暂时不可用 [熔断器状态: HALF_OPEN]
恢复请求  2: 远程服务响应成功 [熔断器状态: HALF_OPEN]
恢复请求  3: 远程服务响应成功 [熔断器状态: CLOSED]
恢复请求  4: 远程服务响应成功 [熔断器状态: CLOSED]
恢复请求  5: 远程服务响应成功 [熔断器状态: CLOSED]

使用Hystrix(比较)

// 使用Hystrix的实现方式(如果有Hystrix依赖)
/*
import com.netflix.hystrix.HystrixCommand;
import com.netflix.hystrix.HystrixCommandGroupKey;
import com.netflix.hystrix.HystrixCommandProperties;
public class HystrixDemo extends HystrixCommand<String> {
    public HystrixDemo() {
        super(Setter.withGroupKey(HystrixCommandGroupKey.Factory.asKey("ExampleGroup"))
            .andCommandPropertiesDefaults(HystrixCommandProperties.Setter()
                .withExecutionTimeoutInMilliseconds(1000)
                .withCircuitBreakerSleepWindowInMilliseconds(5000)
                .withCircuitBreakerErrorThresholdPercentage(50)
                .withCircuitBreakerRequestVolumeThreshold(10)
        ));
    }
    @Override
    protected String run() throws Exception {
        // 实际业务逻辑
        return "success";
    }
    @Override
    protected String getFallback() {
        // 降级逻辑
        return "fallback";
    }
}
*/

这个实现提供了完整的熔断器核心功能:

  1. 三种状态:关闭(CLOSED)、打开(OPEN)、半开(HALF_OPEN)
  2. 失败计数:统计失败次数,达到阈值打开熔断器
  3. 自动恢复:熔断超时后自动进入半开状态
  4. 半开试探:半开状态下限制请求数,成功达到阈值关闭熔断器
  5. 降级处理:熔断时执行降级逻辑
  6. 线程安全:使用原子操作保证并发安全

上一篇Resilience4j案例

下一篇Hystrix案例

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