Resilience4j
概述
Resilience4j 是 Spring Cloud 官方推荐的熔断降级库(替代已废弃的 Netflix Hystrix),提供模块化的弹性能力。与 Hystrix 不同,Resilience4j 使用 函数式编程 和 装饰器模式,核心模块均可独立使用。
六大核心模块
| 模块 | 功能 | 类比 |
|---|---|---|
| CircuitBreaker | 熔断器 | 保险丝(异常率/慢调用率超阈值则断开) |
| RateLimiter | 限流器 | 漏斗(限制每秒请求数) |
| Retry | 重试器 | 自动重试(带退避策略) |
| TimeLimiter | 超时限制 | 闹钟(超时自动取消) |
| Bulkhead | 舱壁隔离 | 隔间(限制并发线程数/信号量) |
| Cache | 缓存 | 结果缓存 |
一、核心模块详解
1.1 CircuitBreaker 熔断
三种状态:
CLOSED (关闭) ── 正常调用
↓ 失败率 ≥ 阈值
OPEN (打开) ── 快速失败
↓ 等待时间过后
HALF_OPEN (半开) ── 尝试放行少量请求
↓ 成功 → CLOSED | 失败 → OPEN配置:
yaml
resilience4j:
circuitbreaker:
configs:
default:
sliding-window-size: 10 # 滑动窗口大小(10 次调用)
minimum-number-of-calls: 5 # 最少调用次数
failure-rate-threshold: 50 # 失败率阈值(百分比)
slow-call-rate-threshold: 50 # 慢调用率阈值
slow-call-duration-threshold: 5s # 慢调用定义(>5s)
wait-duration-in-open-state: 10s # 打开→半开等待时间
permitted-number-of-calls-in-half-open-state: 3 # 半开时允许的请求数
automatic-transition-from-open-to-half-open-enabled: true
record-exceptions:
- java.io.IOException
- java.util.concurrent.TimeoutException
ignore-exceptions:
- org.springframework.web.client.HttpClientErrorException # 4xx 不熔断
payment: # 支付服务独立配置
failure-rate-threshold: 30
wait-duration-in-open-state: 30s
instances:
order-service: # 订单服务使用 default
base-config: default
payment-service: # 支付服务使用 payment
base-config: payment代码使用:
java
@Service
public class PaymentService {
@CircuitBreaker(name = "payment-service", fallbackMethod = "fallback")
public PaymentResult processPayment(PaymentRequest request) {
return paymentClient.charge(request);
}
// fallback 方法:参数签名必须与原始方法一致 + Throwable 参数
public PaymentResult fallback(PaymentRequest request, Throwable t) {
log.error("Payment failed: {}", t.getMessage());
return PaymentResult.failed("Payment temporarily unavailable");
}
}1.2 RateLimiter 限流
yaml
resilience4j:
ratelimiter:
configs:
default:
limit-for-period: 100 # 周期内允许的请求数
limit-refresh-period: 1s # 周期时间
timeout-duration: 500ms # 等待令牌的超时时间
instances:
order-service:
base-config: default
limit-for-period: 200java
@Service
public class OrderService {
@RateLimiter(name = "order-service", fallbackMethod = "rateLimitFallback")
public Order createOrder(CreateOrderRequest request) {
return orderRepository.save(request.toOrder());
}
public Order rateLimitFallback(CreateOrderRequest request, RequestNotPermitted e) {
throw new BusinessException("TOO_MANY_REQUESTS", "Please try again later");
}
}1.3 Retry 重试
yaml
resilience4j:
retry:
configs:
default:
max-attempts: 3 # 最大重试次数
wait-duration: 500ms # 重试间隔
exponential-backoff-multiplier: 2 # 指数退避倍数
retry-exceptions:
- org.springframework.web.client.ResourceAccessException
- java.net.ConnectException
ignore-exceptions:
- org.springframework.web.client.HttpClientErrorException # 4xx 不重试
instances:
logistics-api:
base-config: default
max-attempts: 2java
@Service
public class LogisticsService {
@Retry(name = "logistics-api")
public TrackingResult queryTracking(String orderNo) {
return logisticsClient.query(orderNo); // 网络异常自动重试
}
}1.4 TimeLimiter 超时
java
@Service
public class ExternalApiService {
@TimeLimiter(name = "external-api")
public CompletableFuture<ApiResponse> callExternal() {
return CompletableFuture.supplyAsync(() -> {
// 如果超过 3 秒未返回,抛出 TimeoutException
return restTemplate.getForObject(url, ApiResponse.class);
});
}
}yaml
resilience4j:
timelimiter:
configs:
default:
timeout-duration: 3s # 超时时间
cancel-running-future: true # 超时后取消正在执行的 Future1.5 Bulkhead 舱壁
java
@Service
public class BulkheadService {
@Bulkhead(name = "database-pool", type = Bulkhead.Type.SEMAPHORE,
fallbackMethod = "bulkheadFallback")
public List<User> queryUsers() {
return userRepository.findAll();
}
public List<User> bulkheadFallback(Throwable t) {
log.warn("Bulkhead full, using fallback: {}", t.getMessage());
return List.of(); // 返回空列表
}
}yaml
resilience4j:
bulkhead:
configs:
default:
max-concurrent-calls: 10 # 最大并发数
max-wait-duration: 500ms # 等待超时
instances:
database-pool:
base-config: default
max-concurrent-calls: 5二、Spring Boot 整合
2.1 依赖
xml
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-starter-circuitbreaker-resilience4j</artifactId>
</dependency>
<!-- Actuator 监控(可选) -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
</dependency>2.2 整合 Feign
yaml
spring:
cloud:
openfeign:
circuitbreaker:
enabled: true # 启用 Feign 断路器java
@FeignClient(name = "payment-service", fallback = PaymentClientFallback.class)
public interface PaymentClient {
@PostMapping("/api/payments")
PaymentResult processPayment(PaymentRequest request);
}
@Component
public class PaymentClientFallback implements PaymentClient {
@Override
public PaymentResult processPayment(PaymentRequest request) {
return PaymentResult.failed("Payment service unavailable");
}
}2.3 整合 Gateway
yaml
spring:
cloud:
gateway:
routes:
- id: circuit-breaker-route
uri: lb://unstable-service
predicates:
- Path=/api/unstable/**
filters:
- name: CircuitBreaker
args:
name: gatewayCircuitBreaker
fallbackUri: forward:/fallback三、事件监听与监控
3.1 事件监听
java
@Component
public class Resilience4jEventListener {
@EventListener
public void onCircuitBreakerEvent(CircuitBreakerEvent event) {
switch (event.getEventType()) {
case STATE_TRANSITION:
log.warn("CircuitBreaker state changed: {} → {}",
event.getStateTransition().getFromState(),
event.getStateTransition().getToState());
break;
case FAILURE_RATE_EXCEEDED:
log.error("Failure rate exceeded threshold: {}",
event.getCreationTime());
alertService.sendAlert("CircuitBreaker warning!");
break;
}
}
@EventListener
public void onRetryEvent(RetryEvent event) {
if (event.getEventType() == RetryEvent.Type.RETRY_ON_ERROR) {
log.warn("Retry attempt {} for {}", event.getNumberOfRetryAttempts(),
event.getCreationTime());
}
}
}3.2 指标监控
yaml
management:
endpoints:
web:
exposure:
include: health,metrics,prometheus
metrics:
tags:
application: ${spring.application.name}java
// Actuator 暴露 Resilience4j 指标
// /actuator/metrics/resilience4j.circuitbreaker.calls
// /actuator/metrics/resilience4j.circuitbreaker.state
// /actuator/metrics/resilience4j.ratelimiter.available.permissions
// /actuator/metrics/resilience4j.retry.calls
// 自定义指标
@Component
public class MetricsExporter {
@EventListener
public void onStateTransition(CircuitBreakerEvent event) {
// 发送 Prometheus Counter
if (event.getEventType() == STATE_TRANSITION) {
Metrics.counter("circuitbreaker.state.transition",
"from", event.getStateTransition().getFromState().name(),
"to", event.getStateTransition().getToState().name()
).increment();
}
}
}四、源码架构
Resilience4j 基于 装饰器模式 + 状态机 + 事件发布,核心类围绕 Registry(注册表)、Config(配置)、Module(模块)三层展开。
4.1 核心类层次
java
// Registry:模块注册表(持有配置 + 实例)
public interface CircuitBreakerRegistry {
CircuitBreaker circuitBreaker(String name); // 获取/创建实例
CircuitBreakerConfig getDefaultConfig();
}
// Module:模块接口
public interface CircuitBreaker extends Closeable {
default <T> T executeSupplier(Supplier<T> supplier) { ... }
CircuitBreaker.Metrics getMetrics(); // 统计指标
boolean tryAcquirePermission(); // 尝试获取执行许可
void releasePermission(); // 释放许可
void onError(...); // 记录失败
void onSuccess(...); // 记录成功
}CircuitBreakerRegistry(注册表,单例)
├─ 创建 CircuitBreaker 实例
│ ├─ CircuitBreakerConfig(配置,不可变)
│ ├─ CircuitBreakerMetrics(滑动窗口统计)
│ └─ CircuitBreakerStateMachine(状态机)
└─ 事件处理器(EventProcessor)4.2 CircuitBreaker 状态机源码
java
// 状态机核心:状态转换 + 事件发布
public class CircuitBreakerStateMachine implements CircuitBreaker {
private final AtomicReference<CircuitBreakerState> stateReference;
// 状态转换:CLOSED → OPEN → HALF_OPEN → CLOSED/OPEN
private void transitionToOpenState() {
// 1. 更新状态
stateReference.set(new OpenState(this));
// 2. 发布状态转换事件(观察者收到通知)
publishStateTransitionEvent(CLOSED, OPEN);
}
private void transitionToHalfOpenState() {
stateReference.set(new HalfOpenState(this));
publishStateTransitionEvent(OPEN, HALF_OPEN);
}
// 调用时检查状态
public boolean tryAcquirePermission() {
return stateReference.get().tryAcquirePermission();
}
// 调用结果上报:状态机根据指标决定是否转换
public void onError(long duration, TimeUnit durationUnit, Throwable throwable) {
stateReference.get().onError(duration, durationUnit, throwable);
}
}4.3 各状态实现
java
// CLOSED:正常统计,超过阈值转 OPEN
public class ClosedState extends CircuitBreakerState {
@Override
public void onError(...) {
// 记录失败 → 检查失败率是否超阈值
if (metrics.getFailureRate() >= config.getFailureRateThreshold()) {
stateMachine.transitionToOpenState(); // 触发熔断
}
}
}
// OPEN:直接拒绝
public class OpenState extends CircuitBreakerState {
@Override
public boolean tryAcquirePermission() {
return false; // 熔断期间快速失败
}
// 等待时间到 → HALF_OPEN
}
// HALF_OPEN:放行少量试探请求
public class HalfOpenState extends CircuitBreakerState {
@Override
public boolean tryAcquirePermission() {
return atomicPermission.tryAcquire(); // 只允许限定数量
}
// 成功 → CLOSED;失败 → OPEN
}4.4 滑动窗口统计
java
// 滑动窗口(CountBased 计数窗口 / TimeBased 时间窗口)
public interface CircuitBreakerMetrics {
long getNumberOfSuccessfulCalls(); // 成功次数
long getNumberOfFailedCalls(); // 失败次数
float getFailureRate(); // 失败率
}CountBased:记录最近 N 次调用的结果(环形数组)
TimeBased:记录最近 N 秒内的调用(按时间分桶)4.5 RateLimiter 令牌桶源码
java
public class AtomicRateLimiter implements RateLimiter {
private final AtomicReference<State> state;
// 获取许可
public boolean acquirePermission() {
// 1. 计算可补充的令牌(按时间流逝)
updateStateWithCurrentTime();
// 2. 尝试消费 1 个令牌
return state.get().activePermissions > 0;
}
// 状态快照:令牌数 + 时间戳
private static class State {
final long activePermissions; // 当前可用令牌
final long nanosToWait; // 还需等待的时间
final long cycleStartTime; // 周期开始时间
}
}4.6 事件发布机制
java
// EventProcessor:观察者模式,发布事件
public class EventProcessor<E> implements EventConsumer<E> {
private final List<EventConsumer<E>> consumers = new ArrayList<>();
// 订阅
public void onEvent(EventConsumer<E> consumer) {
consumers.add(consumer);
}
// 发布(状态转换/成功/失败都会触发)
public void publishEvent(E event) {
for (EventConsumer<E> consumer : consumers) {
consumer.consumeEvent(event);
}
}
}java
// 事件类型
public enum CircuitBreakerEvent.Type {
SUCCESS, // 调用成功
ERROR, // 调用失败
STATE_TRANSITION, // 状态转换(CLOSED→OPEN 等)
CALL_NOT_PERMITTED, // 熔断拒绝调用
...
}4.7 Bulkhead 信号量隔离源码
java
public class SemaphoreBulkhead implements Bulkhead {
private final Semaphore semaphore; // 并发许可
@Override
public boolean tryAcquirePermission() {
return semaphore.tryAcquire(); // 并发数超过限制则拒绝
}
@Override
public void releasePermission() {
semaphore.release();
}
}五、配置中心动态调整
5.1 Nacos 动态刷新
yaml
# Nacos 配置中可以动态修改 Resilience4j 参数
resilience4j:
circuitbreaker:
instances:
payment-service:
failure-rate-threshold: 30 # 动态调整:30% → 50%
wait-duration-in-open-state: 30s # 动态调整:30s → 60sjava
// 配置动态刷新后,CircuitBreaker 会自动重新创建
@RefreshScope
@Configuration
public class Resilience4jConfig {
// Resilience4j 自动监听配置变化
}5.2 Java 代码动态调整
java
@Component
public class DynamicConfigAdjuster {
@Autowired
private CircuitBreakerRegistry circuitBreakerRegistry;
@Autowired
private RateLimiterRegistry rateLimiterRegistry;
// 根据流量动态调整限流阈值
public void adjustRateLimiter(String name, int newLimit) {
RateLimiterConfig newConfig = RateLimiterConfig.from(rateLimiterRegistry
.rateLimiter(name).getRateLimiterConfig())
.limitForPeriod(newLimit)
.build();
rateLimiterRegistry.rateLimiter(name).changeTimeoutForDuration(newConfig);
}
// 根据异常率动态调整熔断阈值
public void adjustCircuitBreaker(String name, float failureRate) {
CircuitBreakerConfig newConfig = CircuitBreakerConfig.from(circuitBreakerRegistry
.circuitBreaker(name).getCircuitBreakerConfig())
.failureRateThreshold(failureRate)
.build();
circuitBreakerRegistry.circuitBreaker(name).replaceConfig(newConfig);
}
}六、实战:电商 API 熔断配置
6.1 不同 API 独立熔断策略
yaml
resilience4j:
circuitbreaker:
instances:
payment-gateway: # 支付网关:严格熔断
failure-rate-threshold: 30
wait-duration-in-open-state: 60s
sliding-window-size: 10
sms-service: # 短信服务:宽松熔断
failure-rate-threshold: 60
wait-duration-in-open-state: 5s
sliding-window-size: 20
logistics-query: # 物流查询:快速重试
failure-rate-threshold: 50
wait-duration-in-open-state: 15s
retry:
instances:
logistics-query: # 物流查询:网络波动重试
max-attempts: 3
wait-duration: 200ms
exponential-backoff-multiplier: 2
ratelimiter:
instances:
payment-gateway: # 支付网关:严格限流
limit-for-period: 50
limit-refresh-period: 1s
order-service: # 订单服务:宽松限流
limit-for-period: 500
limit-refresh-period: 1s
bulkhead:
instances:
database-pool: # 数据库连接隔离
max-concurrent-calls: 10
external-api: # 外部 API 隔离
max-concurrent-calls: 56.2 总结
| 模块 | 适用场景 | 关键参数 |
|---|---|---|
| CircuitBreaker | 下游服务异常/慢调用 | sliding-window-size, failure-rate-threshold |
| RateLimiter | 突发流量控制 | limit-for-period, timeout-duration |
| Retry | 网络抖动/临时故障 | max-attempts, exponential-backoff-multiplier |
| TimeLimiter | 第三方 API 超时控制 | timeout-duration |
| Bulkhead | 资源隔离 | max-concurrent-calls |
参考链接: