本文目录导读:

我将为您提供一个完整的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";
}
}
*/
这个实现提供了完整的熔断器核心功能:
- 三种状态:关闭(CLOSED)、打开(OPEN)、半开(HALF_OPEN)
- 失败计数:统计失败次数,达到阈值打开熔断器
- 自动恢复:熔断超时后自动进入半开状态
- 半开试探:半开状态下限制请求数,成功达到阈值关闭熔断器
- 降级处理:熔断时执行降级逻辑
- 线程安全:使用原子操作保证并发安全