一、簡介
斷路器 (斷路器)是微服務架構中的重要設計模式,其靈感來自於現實生活中的斷路器設備。就像斷路器在偵測到過載時自動切斷電路以保護電氣系統一樣,軟體中的斷路器會自動中斷對出現問題的服務的請求以保護整個系統。
這個模式解決什麼問題?
在微服務架構中,服務之間相互依賴。當服務崩潰(緩慢或無回應)時,如果沒有保護機制:
[Service A] ---> [Service B (đang chết)] ---> timeout 30s
↓
Threads bị block
↓
Resource exhaustion
↓
Service A cũng chết (Cascading Failure)
斷路器透過以下方式防止這種骨牌效應 快速失敗 - 立即返回錯誤而不是徒勞等待。
2. 為什麼需要斷路器?
2.1.級聯故障問題
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ Gateway │────▶│ Order Svc │────▶│ Payment Svc │ ← Đang chết
└─────────────┘ └─────────────┘ └─────────────┘
│
▼
┌─────────────┐
│Inventory Svc│
└─────────────┘
當支付服務終止時:
- 訂單服務等待回應→線程被阻塞
- 對 Order Service 的請求增加 → 執行緒池耗盡
- Order Service 無法處理新請求 → 也“死了”
- 網關逾時 → 使用者體驗不佳
2.2.斷路器的好處
| 好處 | 描述 |
|---|---|
| 快速失敗 | 立即返回錯誤,無需等待超時 |
| 保護資源 | 自由線程、連接 |
| 自動恢復 | 服務恢復時自我檢測並恢復 |
| 優雅的降級 | 提供後備響應而不是完整的錯誤 |
| 監控 | 提供有關依賴項運作狀況的指標 |
三、工作原理
3.1.斷路器的三種狀態
┌──────────────────────────────────────┐
│ │
▼ │
┌─────────┐ failure ┌─────────┐ wait timeout ┌─────────────┐
│ CLOSED │───────────────▶│ OPEN │──────────────────▶│ HALF_OPEN │
│(Đóng) │ threshold │ (Mở) │ │ (Nửa mở) │
└─────────┘ └─────────┘ └─────────────┘
▲ ▲ │
│ │ │
│ success │ failure │
└──────────────────────────┴──────────────────────────────┘
CLOSED(關閉)-正常狀態
- 所有請求均被允許透過
- 斷路器監控錯誤率
- 當故障率超過閾值→切換到OPEN
OPEN(開路)-保護狀態
- 所有請求立即被拒絕
- 返回後備響應或異常
- 一段時間後(等待時長)→切換到HALF_OPEN
HALF_OPEN(半開)-測試狀態
- 允許一定數量的請求通過進行測試
- 如果成功→返回CLOSED
- 如果失敗→回傳OPEN
3.2.滑動窗口
Circuit Breaker 使用滑動視窗來計算故障率:
基於計數的滑動視窗:
[Request 1: ✓] [Request 2: ✗] [Request 3: ✓] [Request 4: ✗] [Request 5: ✗]
↑
Failure rate = 3/5 = 60%
基於時間的滑動視窗:
|-------- 10 seconds --------|
| ✓ ✗ ✓ ✗ ✗ ✓ ✗ ✗ ✓ ✓ |
| Failure rate = 5/10 = 50% |
4. 使用 Spring Boot 和 Resilience4j 安裝
4.1.為什麼選擇 Resilience4j?
Hystrix (Netflix) 自 2018 年起已被棄用。 彈性4j 是推薦的函式庫:
- 輕量級,無傳遞依賴
- 專為 Java 8+ 設計,採用函數式編程
- 與 Spring Boot 良好集成
- 響應式支援(Reactor 專案、RxJava)
4.2.依賴關係
<!-- pom.xml --> <dependencyManagement> <dependencies> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-dependencies</artifactId> <version>2023.0.3</version> <type>pom</type> <scope>import</scope> </dependency> </dependencies> </dependencyManagement><dependencies> <!-- Spring Boot Starter --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency>
<!-- Resilience4j Circuit Breaker --> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-starter-circuitbreaker-resilience4j</artifactId> </dependency> <!-- AOP cho annotations --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-aop</artifactId> </dependency> <!-- Actuator cho monitoring --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-actuator</artifactId> </dependency>
</dependencies>
或使用 Gradle:
// build.gradle ext { springCloudVersion = "2023.0.3" }dependencies { implementation 'org.springframework.boot:spring-boot-starter-web' implementation 'org.springframework.cloud:spring-cloud-starter-circuitbreaker-resilience4j' implementation 'org.springframework.boot:spring-boot-starter-aop' implementation 'org.springframework.boot:spring-boot-starter-actuator' }
dependencyManagement { imports { mavenBom "org.springframework.cloud:spring-cloud-dependencies:${springCloudVersion}" } }
5. 詳細配置
5.1.透過 application.yml 配置
# application.yml resilience4j: circuitbreaker: configs: # Cấu hình mặc định cho tất cả circuit breakers default: # Số lượng calls trong sliding window để tính failure rate slidingWindowSize: 10# Loại sliding window: COUNT_BASED hoặc TIME_BASED slidingWindowType: COUNT_BASED # Số calls tối thiểu trước khi tính failure rate minimumNumberOfCalls: 5 # Tỷ lệ lỗi (%) để chuyển sang OPEN failureRateThreshold: 50 # Tỷ lệ slow calls (%) để chuyển sang OPEN slowCallRateThreshold: 100 # Thời gian được coi là slow call (ms) slowCallDurationThreshold: 2000 # Thời gian ở trạng thái OPEN trước khi chuyển sang HALF_OPEN waitDurationInOpenState: 30s # Số calls được phép trong trạng thái HALF_OPEN permittedNumberOfCallsInHalfOpenState: 3 # Tự động chuyển từ OPEN sang HALF_OPEN automaticTransitionFromOpenToHalfOpenEnabled: true # Các exception được ghi nhận là failure recordExceptions: - java.io.IOException - java.net.ConnectException - java.util.concurrent.TimeoutException - org.springframework.web.client.HttpServerErrorException # Các exception KHÔNG được ghi nhận là failure ignoreExceptions: - com.example.BusinessException # Cấu hình cho từng circuit breaker cụ thể instances: # Circuit breaker cho Payment Service paymentService: baseConfig: default failureRateThreshold: 30 waitDurationInOpenState: 20s slidingWindowSize: 20 # Circuit breaker cho Inventory Service inventoryService: baseConfig: default failureRateThreshold: 60 slowCallDurationThreshold: 3000Actuator endpoints
management: endpoints: web: exposure: include: health,circuitbreakers,circuitbreakerevents health: circuitbreakers: enabled: true endpoint: health: show-details: always
5.2.重要參數說明
| 參數 | 描述 | 建議值 |
|---|---|---|
滑動視窗大小 |
計算失敗率的呼叫次數 | 10-100 取決於交通狀況 |
滑動視窗類型 |
COUNT_BASED 或 TIME_BASED | COUNT_BASED(低流量) |
最小呼叫次數 |
評估前最少呼叫次數 | 5-10 |
失敗率閾值 |
開路誤差% | 50%是常見的 |
開啟狀態等待持續時間 |
重試前超時 | 30秒-60秒 |
半開放狀態中允許的呼叫數量 |
HALF_OPEN 中的測試呼叫數量 | 3-10 |
6. 實際例子
6.1.項目結構
order-service/
├── src/main/java/com/example/order/
│ ├── OrderServiceApplication.java
│ ├── config/
│ │ └── Resilience4jConfig.java
│ ├── controller/
│ │ └── OrderController.java
│ ├── service/
│ │ ├── OrderService.java
│ │ └── PaymentServiceClient.java
│ ├── dto/
│ │ ├── OrderRequest.java
│ │ ├── OrderResponse.java
│ │ └── PaymentResponse.java
│ └── exception/
│ └── ServiceUnavailableException.java
├── src/main/resources/
│ └── application.yml
└── pom.xml
6.2.帶有斷路器的支付服務客戶端
package com.example.order.service;import io.github.resilience4j.circuitbreaker.annotation.CircuitBreaker; import io.github.resilience4j.retry.annotation.Retry; import io.github.resilience4j.timelimiter.annotation.TimeLimiter; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import org.springframework.web.client.RestTemplate;
import java.util.concurrent.CompletableFuture;
@Service @RequiredArgsConstructor @Slf4j public class PaymentServiceClient {
private final RestTemplate restTemplate; private static final String PAYMENT_SERVICE_URL = "http://payment-service:8080"; /** * Gọi Payment Service với Circuit Breaker bảo vệ * * Thứ tự xử lý: Retry -> CircuitBreaker -> TimeLimiter -> Bulkhead */ @CircuitBreaker(name = "paymentService", fallbackMethod = "processPaymentFallback") @Retry(name = "paymentService") @TimeLimiter(name = "paymentService") public CompletableFuture<PaymentResponse> processPayment(PaymentRequest request) { log.info("Calling Payment Service for order: {}", request.getOrderId()); return CompletableFuture.supplyAsync(() -> { PaymentResponse response = restTemplate.postForObject( PAYMENT_SERVICE_URL + "/api/payments", request, PaymentResponse.class ); log.info("Payment processed successfully: {}", response); return response; }); } /** * Fallback method khi Circuit Breaker OPEN hoặc có lỗi * * Phải có cùng parameters + thêm Exception/Throwable */ public CompletableFuture<PaymentResponse> processPaymentFallback( PaymentRequest request, Throwable throwable) { log.warn("Payment Service unavailable. Triggering fallback for order: {}. Error: {}", request.getOrderId(), throwable.getMessage()); // Option 1: Trả về response mặc định PaymentResponse fallbackResponse = PaymentResponse.builder() .orderId(request.getOrderId()) .status("PENDING") .message("Payment will be processed when service is available") .fallback(true) .build(); // Option 2: Có thể queue lại để xử lý sau // paymentQueue.add(request); return CompletableFuture.completedFuture(fallbackResponse); } /** * Ví dụ với synchronous call (không dùng TimeLimiter) */ @CircuitBreaker(name = "paymentService", fallbackMethod = "getPaymentStatusFallback") @Retry(name = "paymentService") public PaymentResponse getPaymentStatus(String paymentId) { log.info("Getting payment status for: {}", paymentId); return restTemplate.getForObject( PAYMENT_SERVICE_URL + "/api/payments/" + paymentId, PaymentResponse.class ); } public PaymentResponse getPaymentStatusFallback(String paymentId, Throwable throwable) { log.warn("Cannot get payment status for: {}. Error: {}", paymentId, throwable.getMessage()); return PaymentResponse.builder() .paymentId(paymentId) .status("UNKNOWN") .message("Service temporarily unavailable") .fallback(true) .build(); }
}
6.3.訂單服務
package com.example.order.service;import io.github.resilience4j.circuitbreaker.CircuitBreaker; import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service;
@Service @RequiredArgsConstructor @Slf4j public class OrderService {
private final PaymentServiceClient paymentClient; private final InventoryServiceClient inventoryClient; private final OrderRepository orderRepository; private final CircuitBreakerRegistry circuitBreakerRegistry; public OrderResponse createOrder(OrderRequest request) { log.info("Creating order: {}", request); // 1. Kiểm tra inventory InventoryResponse inventory = inventoryClient.checkInventory(request.getProductId()); if (!inventory.isAvailable()) { throw new InsufficientInventoryException("Product not available"); } // 2. Tạo order với status PENDING Order order = Order.builder() .customerId(request.getCustomerId()) .productId(request.getProductId()) .quantity(request.getQuantity()) .status(OrderStatus.PENDING) .build(); order = orderRepository.save(order); // 3. Xử lý payment (có Circuit Breaker bảo vệ) try { PaymentResponse payment = paymentClient.processPayment( PaymentRequest.builder() .orderId(order.getId()) .amount(request.getAmount()) .build() ).get(); // Blocking call if (payment.isFallback()) { // Payment đang trong queue, cần xử lý sau order.setStatus(OrderStatus.PAYMENT_PENDING); } else if ("SUCCESS".equals(payment.getStatus())) { order.setStatus(OrderStatus.CONFIRMED); } else { order.setStatus(OrderStatus.PAYMENT_FAILED); } } catch (Exception e) { log.error("Payment processing failed", e); order.setStatus(OrderStatus.PAYMENT_PENDING); } order = orderRepository.save(order); return OrderResponse.from(order); } /** * Kiểm tra trạng thái Circuit Breaker programmatically */ public CircuitBreakerStatus getCircuitBreakerStatus(String name) { CircuitBreaker circuitBreaker = circuitBreakerRegistry.circuitBreaker(name); CircuitBreaker.Metrics metrics = circuitBreaker.getMetrics(); return CircuitBreakerStatus.builder() .name(name) .state(circuitBreaker.getState().name()) .failureRate(metrics.getFailureRate()) .slowCallRate(metrics.getSlowCallRate()) .numberOfBufferedCalls(metrics.getNumberOfBufferedCalls()) .numberOfFailedCalls(metrics.getNumberOfFailedCalls()) .numberOfSuccessfulCalls(metrics.getNumberOfSuccessfulCalls()) .numberOfSlowCalls(metrics.getNumberOfSlowCalls()) .build(); }
}
6.4.程式設定(YAML 的替代方案)
package com.example.order.config;import io.github.resilience4j.circuitbreaker.CircuitBreaker; import io.github.resilience4j.circuitbreaker.CircuitBreakerConfig; import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry; import io.github.resilience4j.common.circuitbreaker.configuration.CircuitBreakerConfigCustomizer; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration;
import java.io.IOException; import java.net.ConnectException; import java.time.Duration; import java.util.concurrent.TimeoutException;
@Configuration public class Resilience4jConfig {
@Bean public CircuitBreakerConfigCustomizer paymentServiceCustomizer() { return CircuitBreakerConfigCustomizer.of("paymentService", builder -> builder .slidingWindowType(CircuitBreakerConfig.SlidingWindowType.COUNT_BASED) .slidingWindowSize(10) .minimumNumberOfCalls(5) .failureRateThreshold(50) .slowCallRateThreshold(80) .slowCallDurationThreshold(Duration.ofSeconds(2)) .waitDurationInOpenState(Duration.ofSeconds(30)) .permittedNumberOfCallsInHalfOpenState(3) .automaticTransitionFromOpenToHalfOpenEnabled(true) .recordExceptions( IOException.class, ConnectException.class, TimeoutException.class ) ); } /** * Tạo Circuit Breaker programmatically với custom config */ @Bean public CircuitBreaker customCircuitBreaker() { CircuitBreakerConfig config = CircuitBreakerConfig.custom() .slidingWindowSize(20) .failureRateThreshold(40) .waitDurationInOpenState(Duration.ofSeconds(60)) .permittedNumberOfCallsInHalfOpenState(5) .build(); CircuitBreakerRegistry registry = CircuitBreakerRegistry.of(config); CircuitBreaker circuitBreaker = registry.circuitBreaker("customService"); // Đăng ký event listeners circuitBreaker.getEventPublisher() .onStateTransition(event -> System.out.println("State transition: " + event.getStateTransition())) .onFailureRateExceeded(event -> System.out.println("Failure rate exceeded: " + event.getFailureRate())) .onCallNotPermitted(event -> System.out.println("Call not permitted")) .onError(event -> System.out.println("Error: " + event.getThrowable().getMessage())); return circuitBreaker; }
}
6.5.與 WebClient 一起使用(反應式)
package com.example.order.service;import io.github.resilience4j.circuitbreaker.annotation.CircuitBreaker; import io.github.resilience4j.reactor.circuitbreaker.operator.CircuitBreakerOperator; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import org.springframework.web.reactive.function.client.WebClient; import reactor.core.publisher.Mono;
@Service @RequiredArgsConstructor @Slf4j public class PaymentServiceReactiveClient {
private final WebClient webClient; private final io.github.resilience4j.circuitbreaker.CircuitBreaker circuitBreaker; /** * Cách 1: Sử dụng annotation */ @CircuitBreaker(name = "paymentService", fallbackMethod = "processPaymentFallback") public Mono<PaymentResponse> processPayment(PaymentRequest request) { return webClient.post() .uri("/api/payments") .bodyValue(request) .retrieve() .bodyToMono(PaymentResponse.class) .doOnSuccess(response -> log.info("Payment successful: {}", response)) .doOnError(error -> log.error("Payment failed: {}", error.getMessage())); } public Mono<PaymentResponse> processPaymentFallback(PaymentRequest request, Throwable t) { log.warn("Fallback triggered for payment: {}", request.getOrderId()); return Mono.just(PaymentResponse.builder() .orderId(request.getOrderId()) .status("PENDING") .fallback(true) .build()); } /** * Cách 2: Sử dụng operator programmatically */ public Mono<PaymentResponse> processPaymentWithOperator(PaymentRequest request) { return webClient.post() .uri("/api/payments") .bodyValue(request) .retrieve() .bodyToMono(PaymentResponse.class) .transformDeferred(CircuitBreakerOperator.of(circuitBreaker)) .onErrorResume(throwable -> { log.warn("Circuit breaker triggered fallback"); return Mono.just(PaymentResponse.builder() .status("PENDING") .fallback(true) .build()); }); }
}
6.6。具有斷路器狀態的 REST 控制器
package com.example.order.controller;import io.github.resilience4j.circuitbreaker.CircuitBreaker; import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry; import lombok.RequiredArgsConstructor; import org.springframework.http.ResponseEntity; import org.springframework.web.bind.annotation.*;
import java.util.HashMap; import java.util.Map;
@RestController @RequestMapping("/api/orders") @RequiredArgsConstructor public class OrderController {
private final OrderService orderService; private final CircuitBreakerRegistry circuitBreakerRegistry; @PostMapping public ResponseEntity<OrderResponse> createOrder(@RequestBody OrderRequest request) { OrderResponse response = orderService.createOrder(request); return ResponseEntity.ok(response); } /** * Endpoint kiểm tra trạng thái Circuit Breakers */ @GetMapping("/circuit-breakers/status") public ResponseEntity<Map<String, Object>> getCircuitBreakersStatus() { Map<String, Object> status = new HashMap<>(); circuitBreakerRegistry.getAllCircuitBreakers().forEach(cb -> { CircuitBreaker.Metrics metrics = cb.getMetrics(); Map<String, Object> cbStatus = new HashMap<>(); cbStatus.put("state", cb.getState().name()); cbStatus.put("failureRate", metrics.getFailureRate()); cbStatus.put("slowCallRate", metrics.getSlowCallRate()); cbStatus.put("bufferedCalls", metrics.getNumberOfBufferedCalls()); cbStatus.put("failedCalls", metrics.getNumberOfFailedCalls()); cbStatus.put("successfulCalls", metrics.getNumberOfSuccessfulCalls()); cbStatus.put("notPermittedCalls", metrics.getNumberOfNotPermittedCalls()); status.put(cb.getName(), cbStatus); }); return ResponseEntity.ok(status); } /** * Endpoint để manually transition Circuit Breaker * (Chỉ dùng cho testing/debugging) */ @PostMapping("/circuit-breakers/{name}/transition/{state}") public ResponseEntity<String> transitionCircuitBreaker( @PathVariable String name, @PathVariable String state) { CircuitBreaker circuitBreaker = circuitBreakerRegistry.circuitBreaker(name); switch (state.toUpperCase()) { case "OPEN": circuitBreaker.transitionToOpenState(); break; case "CLOSED": circuitBreaker.transitionToClosedState(); break; case "HALF_OPEN": circuitBreaker.transitionToHalfOpenState(); break; default: return ResponseEntity.badRequest().body("Invalid state"); } return ResponseEntity.ok("Transitioned to " + state); }
}
7. 什麼時候應該使用斷路器?
7.1.在下列情況下應使用斷路器:
✅ 致電外部服務
// Gọi Payment Gateway (Stripe, VNPay, MoMo) @CircuitBreaker(name = "paymentGateway") public PaymentResult processPayment(PaymentRequest request) { return paymentGatewayClient.charge(request); }
// Gọi SMS/Email Provider @CircuitBreaker(name = "notificationService") public void sendNotification(NotificationRequest request) { twilioClient.sendSMS(request); }
✅ 微服務之間的通信
// Order Service → Inventory Service @CircuitBreaker(name = "inventoryService") public InventoryStatus checkStock(String productId) { return inventoryClient.getStock(productId); }
// User Service → Auth Service @CircuitBreaker(name = "authService") public TokenInfo validateToken(String token) { return authClient.validate(token); }
✅ 透過網路存取資料庫
// Database cluster có thể unavailable
@CircuitBreaker(name = "databaseOperation")
public List<Order> getOrderHistory(String customerId) {
return orderRepository.findByCustomerId(customerId);
}
✅ 呼叫第三方API
// Weather API, Exchange Rate API, Social Login
@CircuitBreaker(name = "weatherApi")
public WeatherData getCurrentWeather(String location) {
return weatherApiClient.fetch(location);
}
7.2.在下列情況下不應使用斷路器:
❌ 內部操作沒有網路調用
// Tính toán local - KHÔNG cần Circuit Breaker public BigDecimal calculateDiscount(Order order) { return order.getTotal().multiply(DISCOUNT_RATE); }
// Validate input - KHÔNG cần Circuit Breaker public boolean validateEmail(String email) { return EMAIL_PATTERN.matcher(email).matches(); }
❌ 逾時和重試就夠了
// Nếu đã có cơ chế retry với exponential backoff
// và timeout hợp lý, có thể không cần thêm Circuit Breaker
❌ 關鍵操作不能失敗
// Ví dụ: Ghi log audit bắt buộc
// Không nên dùng Circuit Breaker vì không thể skip
7.3.決策矩陣
| 情況 | 斷路器? | 原因 |
|---|---|---|
| 來電支付服務 | ✅ 是的 | 外部服務,可下載 |
| 呼叫內部緩存 | ❌ 沒有 | 本地,無網路延遲 |
| 呼叫資料庫主資料庫 | ⚠️考慮 | 根據設置,通常有一個連接池 |
| 呼叫訊息隊列 | ✅ 是的 | 網路調用,可能超時 |
| 讀取本地文件 | ❌ 沒有 | 本地 I/O,不需要 |
| 呼叫 OAuth 提供者 | ✅ 是的 | 外部服務,關鍵 |
| CPU密集型運算 | ❌ 沒有 | 不是網路問題 |
7.4.與其他模式結合
# Thứ tự áp dụng (từ ngoài vào trong): # Retry -> CircuitBreaker -> RateLimiter -> TimeLimiter -> Bulkheadresilience4j: retry: instances: paymentService: maxAttempts: 3 waitDuration: 1s retryExceptions: - java.io.IOException
circuitbreaker: instances: paymentService: failureRateThreshold: 50 waitDurationInOpenState: 30s
ratelimiter: instances: paymentService: limitForPeriod: 100 limitRefreshPeriod: 1s
timelimiter: instances: paymentService: timeoutDuration: 5s
bulkhead: instances: paymentService: maxConcurrentCalls: 10 maxWaitDuration: 500ms
8. 監控和可觀察性
8.1.執行器端點
# application.yml
management:
endpoints:
web:
exposure:
include: health,circuitbreakers,circuitbreakerevents
health:
circuitbreakers:
enabled: true
endpoint:
health:
show-details: always
可用端點:
# Xem tất cả circuit breakers GET /actuator/circuitbreakersXem events của circuit breaker cụ thể
GET /actuator/circuitbreakerevents GET /actuator/circuitbreakerevents/{name}
Health check bao gồm circuit breaker status
GET /actuator/health
回應範例:
// GET /actuator/circuitbreakers
{
"circuitBreakers": {
"paymentService": {
"failureRate": "25.0%",
"slowCallRate": "10.0%",
"failureRateThreshold": "50.0%",
"slowCallRateThreshold": "100.0%",
"bufferedCalls": 20,
"failedCalls": 5,
"slowCalls": 2,
"slowFailedCalls": 1,
"notPermittedCalls": 0,
"state": "CLOSED"
}
}
}
8.2. Prometheus + Grafana 集成
<!-- Thêm dependency -->
<dependency>
<groupId>io.micrometer</groupId>
<artifactId>micrometer-registry-prometheus</artifactId>
</dependency>
# application.yml
management:
endpoints:
web:
exposure:
include: health,prometheus,circuitbreakers
prometheus:
metrics:
export:
enabled: true
可用指標:
# Số lượng calls theo trạng thái resilience4j_circuitbreaker_calls_total{name="paymentService",kind="successful"} 150 resilience4j_circuitbreaker_calls_total{name="paymentService",kind="failed"} 10 resilience4j_circuitbreaker_calls_total{name="paymentService",kind="not_permitted"} 5Trạng thái circuit breaker (0=CLOSED, 1=OPEN, 2=HALF_OPEN)
resilience4j_circuitbreaker_state{name="paymentService",state="closed"} 1
Failure rate
resilience4j_circuitbreaker_failure_rate{name="paymentService"} 6.25
Slow call rate
resilience4j_circuitbreaker_slow_call_rate{name="paymentService"} 2.5
8.3.自訂事件處理程序
package com.example.order.config;import io.github.resilience4j.circuitbreaker.CircuitBreaker; import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry; import io.github.resilience4j.circuitbreaker.event.*; import lombok.extern.slf4j.Slf4j; import org.springframework.context.annotation.Configuration;
import jakarta.annotation.PostConstruct;
@Configuration @Slf4j public class CircuitBreakerEventConfig {
private final CircuitBreakerRegistry registry; private final AlertService alertService; public CircuitBreakerEventConfig(CircuitBreakerRegistry registry, AlertService alertService) { this.registry = registry; this.alertService = alertService; } @PostConstruct public void registerEventConsumers() { registry.getAllCircuitBreakers().forEach(this::registerEventConsumer); // Đăng ký cho các circuit breakers được tạo sau registry.getEventPublisher() .onEntryAdded(event -> registerEventConsumer(event.getAddedEntry())); } private void registerEventConsumer(CircuitBreaker circuitBreaker) { circuitBreaker.getEventPublisher() .onStateTransition(this::handleStateTransition) .onFailureRateExceeded(this::handleFailureRateExceeded) .onCallNotPermitted(this::handleCallNotPermitted) .onError(this::handleError) .onSuccess(this::handleSuccess); } private void handleStateTransition(CircuitBreakerOnStateTransitionEvent event) { String message = String.format( "Circuit Breaker '%s' transitioned from %s to %s", event.getCircuitBreakerName(), event.getStateTransition().getFromState(), event.getStateTransition().getToState() ); log.warn(message); // Gửi alert khi chuyển sang OPEN if (event.getStateTransition().getToState() == CircuitBreaker.State.OPEN) { alertService.sendAlert( AlertLevel.HIGH, "Circuit Breaker OPENED", message ); } // Gửi notification khi recovery (về CLOSED) if (event.getStateTransition().getToState() == CircuitBreaker.State.CLOSED && event.getStateTransition().getFromState() == CircuitBreaker.State.HALF_OPEN) { alertService.sendNotification( "Circuit Breaker RECOVERED", message ); } } private void handleFailureRateExceeded(CircuitBreakerOnFailureRateExceededEvent event) { log.error("Failure rate exceeded for '{}': {}%", event.getCircuitBreakerName(), event.getFailureRate()); } private void handleCallNotPermitted(CircuitBreakerOnCallNotPermittedEvent event) { log.debug("Call not permitted for '{}'", event.getCircuitBreakerName()); } private void handleError(CircuitBreakerOnErrorEvent event) { log.error("Error in '{}': {}", event.getCircuitBreakerName(), event.getThrowable().getMessage()); } private void handleSuccess(CircuitBreakerOnSuccessEvent event) { log.trace("Successful call to '{}' in {}ms", event.getCircuitBreakerName(), event.getElapsedDuration().toMillis()); }
}
9. 最佳實踐
9.1.後備策略設計
/**
Các chiến lược Fallback phổ biến */ public class FallbackStrategies {
// 1. Trả về giá trị mặc định public ProductInfo getProductFallback(String productId, Throwable t) { return ProductInfo.builder() .id(productId) .name("Product information temporarily unavailable") .available(false) .build(); }
// 2. Trả về dữ liệu từ cache @Autowired private CacheManager cacheManager;
public ProductInfo getProductFromCacheFallback(String productId, Throwable t) { Cache cache = cacheManager.getCache("products"); ProductInfo cached = cache.get(productId, ProductInfo.class);
if (cached != null) { cached.setFromCache(true); return cached; } return getProductFallback(productId, t);}
// 3. Gọi service backup public ProductInfo getProductFromBackupServiceFallback(String productId, Throwable t) { try { return backupProductService.getProduct(productId); } catch (Exception e) { return getProductFromCacheFallback(productId, t); } }
// 4. Queue để xử lý sau (async fallback) @Autowired private RabbitTemplate rabbitTemplate;
public OrderResponse createOrderAsyncFallback(OrderRequest request, Throwable t) { // Lưu vào queue để xử lý sau rabbitTemplate.convertAndSend("order-retry-queue", request);
return OrderResponse.builder() .status("QUEUED") .message("Your order is being processed") .estimatedProcessingTime("5 minutes") .build();}
// 5. Graceful degradation - giảm tính năng public RecommendationResponse getRecommendationsFallback(String userId, Throwable t) { // Thay vì personalized recommendations, trả về popular items return RecommendationResponse.builder() .items(popularItemsCache.getTopItems(10)) .type("POPULAR") // Thay vì "PERSONALIZED" .message("Showing popular items") .build(); } }
9.2.調整參數
# Development/Testing environment resilience4j: circuitbreaker: configs: development: slidingWindowSize: 5 # Window nhỏ để test nhanh minimumNumberOfCalls: 3 # Ít calls để trigger sớm failureRateThreshold: 50 waitDurationInOpenState: 10s # Recovery nhanh permittedNumberOfCallsInHalfOpenState: 2Production environment
resilience4j: circuitbreaker: configs: production: slidingWindowSize: 100 # Window lớn hơn, ổn định hơn minimumNumberOfCalls: 20 # Cần nhiều data points failureRateThreshold: 50 slowCallRateThreshold: 80 # Theo dõi slow calls slowCallDurationThreshold: 3s waitDurationInOpenState: 60s # Chờ lâu hơn permittedNumberOfCallsInHalfOpenState: 10
9.3.測試斷路器
package com.example.order.service;import io.github.resilience4j.circuitbreaker.CircuitBreaker; import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.boot.test.mock.mockito.MockBean;
import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.*;
@SpringBootTest class PaymentServiceClientTest {
@Autowired private PaymentServiceClient paymentClient; @Autowired private CircuitBreakerRegistry circuitBreakerRegistry; @MockBean private RestTemplate restTemplate; private CircuitBreaker circuitBreaker; @BeforeEach void setUp() { circuitBreaker = circuitBreakerRegistry.circuitBreaker("paymentService"); circuitBreaker.reset(); // Reset state before each test } @Test void shouldReturnSuccessfulResponse() { // Given PaymentResponse expectedResponse = PaymentResponse.builder() .status("SUCCESS") .build(); when(restTemplate.postForObject(any(), any(), any())) .thenReturn(expectedResponse); // When PaymentResponse response = paymentClient.processPayment( PaymentRequest.builder().orderId("123").build() ).join(); // Then assertThat(response.getStatus()).isEqualTo("SUCCESS"); assertThat(response.isFallback()).isFalse(); assertThat(circuitBreaker.getState()).isEqualTo(CircuitBreaker.State.CLOSED); } @Test void shouldOpenCircuitAfterFailures() { // Given - Circuit breaker có threshold 50%, window size 10 when(restTemplate.postForObject(any(), any(), any())) .thenThrow(new RuntimeException("Service unavailable")); // When - Gọi nhiều lần để trigger circuit breaker for (int i = 0; i < 10; i++) { try { paymentClient.processPayment( PaymentRequest.builder().orderId("123").build() ).join(); } catch (Exception ignored) { } } // Then assertThat(circuitBreaker.getState()).isEqualTo(CircuitBreaker.State.OPEN); } @Test void shouldReturnFallbackWhenCircuitOpen() { // Given circuitBreaker.transitionToOpenState(); // When PaymentResponse response = paymentClient.processPayment( PaymentRequest.builder().orderId("123").build() ).join(); // Then assertThat(response.isFallback()).isTrue(); assertThat(response.getStatus()).isEqualTo("PENDING"); verify(restTemplate, never()).postForObject(any(), any(), any()); } @Test void shouldTransitionToHalfOpenAfterWaitDuration() throws InterruptedException { // Given circuitBreaker.transitionToOpenState(); // When - Chờ hết waitDurationInOpenState (đã config là 10s trong test profile) Thread.sleep(11000); // Then assertThat(circuitBreaker.getState()).isEqualTo(CircuitBreaker.State.HALF_OPEN); } @Test void shouldCloseAfterSuccessfulCallsInHalfOpen() { // Given circuitBreaker.transitionToHalfOpenState(); when(restTemplate.postForObject(any(), any(), any())) .thenReturn(PaymentResponse.builder().status("SUCCESS").build()); // When - Số calls thành công >= permittedNumberOfCallsInHalfOpenState for (int i = 0; i < 3; i++) { paymentClient.processPayment( PaymentRequest.builder().orderId("123").build() ).join(); } // Then assertThat(circuitBreaker.getState()).isEqualTo(CircuitBreaker.State.CLOSED); }
}
9.4.要避免的常見錯誤
/**
- ❌ SAI: Fallback method signature không đúng */ @CircuitBreaker(name = "service", fallbackMethod = "fallback") public String callService(String param) { return service.call(param); }
// ❌ Thiếu Throwable parameter public String fallback(String param) { // Sẽ không được gọi! return "default"; }
// ✅ ĐÚNG: Phải có Throwable/Exception public String fallback(String param, Throwable t) { return "default"; }
/**
- ❌ SAI: Catch exception trong method */ @CircuitBreaker(name = "service") public String callService() { try { return service.call(); } catch (Exception e) { return "error"; // Circuit Breaker không thấy failure! } }
// ✅ ĐÚNG: Để exception propagate @CircuitBreaker(name = "service", fallbackMethod = "fallback") public String callService() { return service.call(); // Throw exception nếu fail }
/**
❌ SAI: Self-invocation (gọi method trong cùng class) */ @Service public class MyService {
@CircuitBreaker(name = "service") public String callService() { return externalService.call(); }
public String doSomething() { return callService(); // Circuit Breaker không work! } }
// ✅ ĐÚNG: Gọi từ class khác hoặc inject self @Service public class MyService {
@Autowired private MyService self; // Hoặc inject từ ApplicationContext @CircuitBreaker(name = "service") public String callService() { return externalService.call(); } public String doSomething() { return self.callService(); // Đi qua proxy, Circuit Breaker work! }}
/**
- ❌ SAI: Không set timeout, leading to thread exhaustion */ @CircuitBreaker(name = "service") public String callSlowService() { return slowService.call(); // Có thể block 60s! }
// ✅ ĐÚNG: Kết hợp với TimeLimiter @CircuitBreaker(name = "service") @TimeLimiter(name = "service") // Timeout sau 5s public CompletableFuture<String> callSlowService() { return CompletableFuture.supplyAsync(() -> slowService.call()); }
10. 結論
10.1.總結
斷路器是微服務架構中不可或缺的模式,可以:
- 防止級聯故障 - 保護您的系統免受骨牌效應的影響
- 快速失敗 - 立即返回錯誤而不是等待超時
- 自我修復 - 當依賴服務再次工作時自動恢復
- 優雅的降級 - 提供後備方案而不是完全失敗
10.2.部署清單
- [ ] 識別需要保護的外部依賴項
- [ ] 新增 Resilience4j 依賴項
- [ ] 適當配置斷路器參數
- [ ] 為所有斷路器實施回退方法
- [ ] 使用Actuator + Prometheus添加監控
- [ ] 設定狀態轉換警報
- [ ] 為斷路器行為編寫單元測試
- [ ] 根據生產指標調整參數
10.3.資源
本文是為 Spring Boot 高級系列編寫的 - 適合想要建立彈性微服務系統的後端開發人員和 DevOps 工程師。
