Chuyển đến nội dung chính

レッスン 16: フォールト トレランス — サーキット ブレーカー、再試行、バルクヘッド

SmallRye フォールト トレランス: @Retry、@Timeout、@CircuitBreaker、@Bulkhead、@Fallback — 連鎖的な障害からサービスを保護します。

💻 プログラミング — レッスン 15 レッスン 16: フォールト トレランス — サーキット ブレーカー 再試行、バルクヘッド

Quarkus マイクロサービス: 基本から運用まで

パート 5: 回復力と可観測性

xdev.asia

はじめに

マイクロサービスでは、サービスが遅くなったりダウンしたりすると、システム全体に影響を与える可能性があります (カスケード障害)。 SmallRye フォールト トレランス (MicroProfile Fault Tolerance) は、再試行、タイムアウト、サーキット ブレーカー、バルクヘッド、フォールバックといった保護パターンを提供します。

依存関係

<dependency>
    <groupId>io.quarkus</groupId>
    <artifactId>quarkus-smallrye-fault-tolerance</artifactId>
</dependency>

@Retry — 自動的に再試行します

import org.eclipse.microprofile.faulttolerance.Retry;

@ApplicationScoped
public class OrderService {

    @Inject @RestClient
    ProductServiceClient productClient;

    @Retry(maxRetries = 3,
           delay = 500,           // 500ms giữa các lần retry
           jitter = 200,          // ±200ms random delay
           retryOn = {IOException.class,
                      WebApplicationException.class},
           abortOn = {ResourceNotFoundException.class})
    public ProductInfo getProduct(Long productId) {
        return productClient.getById(productId);
    }
}

指数バックオフ

@Retry(maxRetries = 4,
       delay = 1000,
       maxDuration = 30000)
@ExponentialBackoff(factor = 2, maxDelay = 10000)
// Retry: 1s → 2s → 4s → 8s (max 10s)
public ProductInfo getProductWithBackoff(Long id) {
    return productClient.getById(id);
}

@Timeout — 時間制限

import org.eclipse.microprofile.faulttolerance.Timeout;

@Timeout(value = 5, unit = ChronoUnit.SECONDS)
public ProductInfo getProduct(Long id) {
    // Nếu > 5s → throw TimeoutException
    return productClient.getById(id);
}

再試行とタイムアウトを組み合わせる

@Retry(maxRetries = 3, delay = 1000)
@Timeout(5000)  // Mỗi lần try tối đa 5s
public ProductInfo getProductReliable(Long id) {
    return productClient.getById(id);
}
// Worst case: 3 retries × 5s timeout + 3 × 1s delay = 18s max

@CircuitBreaker — サーキットブレーカー

サーキット ブレーカーは、障害が発生したサービスへの呼び出しを「中断」することで、連鎖的な障害を防ぎます。

     ┌─────────┐    failures > threshold    ┌──────┐
     │ CLOSED  │ ──────────────────────────> │ OPEN │
     │(normal) │                             │(fail)│
     └─────────┘                             └──┬───┘
          ^                                     │
          │        ┌────────────┐               │
          │        │ HALF-OPEN  │ <─────────────┘
          └────────│ (testing)  │  after delay
    success        └────────────┘
import org.eclipse.microprofile.faulttolerance.CircuitBreaker;

@CircuitBreaker(
    requestVolumeThreshold = 20,  // Window: 20 requests
    failureRatio = 0.5,           // Mở khi 50% fail
    delay = 10,                   // Đợi 10s trước khi thử lại
    delayUnit = ChronoUnit.SECONDS,
    successThreshold = 3)         // Đóng lại sau 3 success
@Fallback(fallbackMethod = "getProductFallback")
public ProductInfo getProduct(Long id) {
    return productClient.getById(id);
}

// Fallback khi circuit OPEN hoặc call fail
public ProductInfo getProductFallback(Long id) {
    // Trả về cached data hoặc default
    return cachedProducts.getOrDefault(id,
        new ProductInfo(id, "Product Unavailable",
            BigDecimal.ZERO, "VND"));
}

@Bulkhead — 同時呼び出しを制限する

同時リクエストが多すぎることによってサービスが過負荷にならないようにします。

import org.eclipse.microprofile.faulttolerance.Bulkhead;

@Bulkhead(value = 10,             // Max 10 concurrent calls
          waitingTaskQueue = 5)    // Max 5 queued
@Timeout(5000)
public ProductInfo getProduct(Long id) {
    return productClient.getById(id);
}
// Nếu > 15 (10 running + 5 queued) → BulkheadException

@Fallback — 値を置換する

import org.eclipse.microprofile.faulttolerance.Fallback;

// Method-level fallback
@Fallback(fallbackMethod = "getProductFallback")
@Retry(maxRetries = 2)
@Timeout(3000)
public ProductInfo getProduct(Long id) {
    return productClient.getById(id);
}

private ProductInfo getProductFallback(Long id) {
    Log.warnf("Fallback for product %d", id);
    // Trả về từ local cache
    return productCache.get(id);
}

// Handler class fallback
@Fallback(value = ProductFallbackHandler.class)
public ProductInfo getProduct2(Long id) {
    return productClient.getById(id);
}

public class ProductFallbackHandler
        implements FallbackHandler<ProductInfo> {
    @Override
    public ProductInfo handle(ExecutionContext context) {
        // Log error, return default
        return new ProductInfo(
            0L, "Unavailable", BigDecimal.ZERO, "VND");
    }
}

アノテーションを組み合わせる場合の実行順序

複数のアノテーションを同時に使用する場合、実行順序は次のとおりです。

Request → Bulkhead → CircuitBreaker → Retry → Timeout → Method → Fallback
                                                                      ↑
                                                              (on any failure)
Ví dụ thực tế với getProduct():
┌─────────────────────────────────────────────────────────┐
│ ① Bulkhead: Có slot trống?                              │
│   ├─ YES → tiếp tục                                    │
│   └─ NO  → BulkheadException → ⑥ Fallback              │
│                                                         │
│ ② CircuitBreaker: CLOSED?                               │
│   ├─ CLOSED → tiếp tục                                 │
│   ├─ HALF-OPEN → cho 1 request thử                     │
│   └─ OPEN  → CircuitBreakerOpenException → ⑥ Fallback  │
│                                                         │
│ ③ Retry: lần thử thứ mấy? (max 3)                      │
│   ├─ Lần 1 → gọi method                                │
│   └─ Fail → đợi delay → retry lần 2, 3...              │
│                                                         │
│ ④ Timeout: method chạy < 5s?                            │
│   ├─ YES → trả kết quả                                 │
│   └─ NO  → TimeoutException → ③ Retry thử lại          │
│                                                         │
│ ⑤ Method: productClient.getById(id)                     │
│                                                         │
│ ⑥ Fallback: khi tất cả retries fail                     │
│   → getProductFallback(id)                              │
└─────────────────────────────────────────────────────────┘

すべてを組み合わせる — 現実世界のパターン

@ApplicationScoped
public class ResilientProductClient {

    @Inject @RestClient
    ProductServiceClient productClient;

    @Inject
    ProductCache productCache;

    @Retry(maxRetries = 3, delay = 500, jitter = 200,
           retryOn = IOException.class)
    @CircuitBreaker(requestVolumeThreshold = 20,
                    failureRatio = 0.5,
                    delay = 10, delayUnit = ChronoUnit.SECONDS)
    @Bulkhead(value = 20, waitingTaskQueue = 10)
    @Timeout(5000)
    @Fallback(fallbackMethod = "getProductFallback")
    public ProductInfo getProduct(Long id) {
        ProductInfo product = productClient.getById(id);
        // Cập nhật cache khi thành công
        productCache.put(id, product);
        return product;
    }

    public ProductInfo getProductFallback(Long id) {
        Log.warnf("Using cached product for id: %d", id);
        ProductInfo cached = productCache.get(id);
        if (cached != null) return cached;

        throw new ServiceUnavailableException(
            "Product Service unavailable "
            + "and no cached data for product " + id);
    }
}

@RateLimit — リクエストレート制限 (Quarkus 3.x)

import io.smallrye.faulttolerance.api.RateLimit;
import java.time.temporal.ChronoUnit;

// Giới hạn 100 requests / 1 phút
@RateLimit(value = 100,
           window = 1, windowUnit = ChronoUnit.MINUTES)
@Fallback(fallbackMethod = "rateLimitedFallback")
public ProductInfo getProduct(Long id) {
    return productClient.getById(id);
}

public ProductInfo rateLimitedFallback(Long id) {
    throw new WebApplicationException(
        "Rate limit exceeded. Try again later.", 429);
}

連鎖的な障害 — 実際の例

Payment Service が遅い (DB 過負荷) と仮定します。

┌──────────┐     ┌──────────┐     ┌──────────┐
│  Client  │────>│  Order   │────>│ Payment  │ ← DB chậm (30s response)
│ Browser  │     │ Service  │     │ Service  │
└──────────┘     └──────────┘     └──────────┘
     │                │
     │  Thread pool   │ Tất cả threads blocked
     │  cạn kiệt!     │ chờ Payment response
     │                ▼
     │           Order Service
     │           KHÔNG THỂ xử lý
     │           request mới
     └────> 503 Service Unavailable

フォールト トレランスなし: 支払いサービスとともに注文サービスが停止します (カスケード障害)。

フォールト トレランスを備えています:

@ApplicationScoped
public class ResilientPaymentClient {

    @Inject @RestClient
    PaymentServiceClient paymentClient;

    @Timeout(3000)          // Không chờ quá 3s
    @CircuitBreaker(
        requestVolumeThreshold = 10,
        failureRatio = 0.5,
        delay = 30,
        delayUnit = ChronoUnit.SECONDS)
    @Bulkhead(value = 5)   // Max 5 concurrent calls
    @Fallback(fallbackMethod = "paymentFallback")
    public PaymentResult processPayment(
            PaymentRequest request) {
        return paymentClient.process(request);
    }

    public PaymentResult paymentFallback(
            PaymentRequest request) {
        Log.warnf("Payment Service unavailable, "
            + "queuing payment for order %s",
            request.orderId());

        // Gửi vào Kafka queue để retry sau
        paymentQueue.send(request);

        return new PaymentResult(
            request.orderId(),
            "PENDING",
            "Payment queued for processing");
    }
}

結果:

  • タイムアウト 3 秒 → 30 秒間スレッドをブロックしません
  • サーキット ブレーカー → 5/10 回の障害の後、すぐにブレーク → フェイルファスト
  • バルクヘッド 5 → 5 つのスレッドのみが Payment を呼び出し、残りは他のリクエストを処理します
  • フォールバック → Kafka にキューに入り、バックグラウンドで再試行 → ユーザーはブロックされない

メトリクスの統合

SmallRye Fault Tolerance は、利用可能な場合にはメトリクスを自動的に公開します quarkus-micrometer:

<dependency>
    <groupId>io.quarkus</groupId>
    <artifactId>quarkus-micrometer-registry-prometheus</artifactId>
</dependency>

自動メトリクスは次の場所から入手できます。 /q/metrics:

# Retry metrics
ft_retry_calls_total{method="getProduct",retried="true",retryResult="valueReturned"} 42
ft_retry_calls_total{method="getProduct",retried="true",retryResult="maxRetriesReached"} 3
ft_retry_retries_total{method="getProduct"} 86

# Circuit Breaker metrics
ft_circuitbreaker_state_total{method="getProduct",state="closed"} 95.0
ft_circuitbreaker_state_total{method="getProduct",state="open"} 3.0
ft_circuitbreaker_state_total{method="getProduct",state="halfOpen"} 2.0
ft_circuitbreaker_opened_total{method="getProduct"} 2

# Bulkhead metrics
ft_bulkhead_executionsRunning{method="getProduct"} 3
ft_bulkhead_executionsWaiting{method="getProduct"} 1
ft_bulkhead_runningDuration_seconds{method="getProduct",quantile="0.95"} 0.234

# Timeout metrics
ft_timeout_calls_total{method="getProduct",timedOut="true"} 5
ft_timeout_calls_total{method="getProduct",timedOut="false"} 195

Grafana ダッシュボード クエリの例

# Circuit Breaker open rate
rate(ft_circuitbreaker_opened_total[5m])

# Retry success rate
rate(ft_retry_calls_total{retryResult="valueReturned"}[5m])
/ rate(ft_retry_calls_total[5m])

# P95 bulkhead execution time
ft_bulkhead_runningDuration_seconds{quantile="0.95"}

application.properties による設定

# Override annotations qua config
com.xdev.ecommerce.order.ResilientProductClient/getProduct/Retry/maxRetries=5
com.xdev.ecommerce.order.ResilientProductClient/getProduct/Timeout/value=3000
com.xdev.ecommerce.order.ResilientProductClient/getProduct/CircuitBreaker/delay=15000

# Global defaults
Retry/maxRetries=3
Timeout/value=5000
CircuitBreaker/failureRatio=0.5

フォールト トレランスのテスト

@QuarkusTest
class ResilientProductClientTest {

    @InjectMock
    @RestClient
    ProductServiceClient productClient;

    @Inject
    ResilientProductClient resilientClient;

    @Test
    void testRetryOnFailure() {
        // First two calls fail, third succeeds
        when(productClient.getById(1L))
            .thenThrow(new IOException("timeout"))
            .thenThrow(new IOException("timeout"))
            .thenReturn(new ProductInfo(1L, "Test",
                BigDecimal.TEN, "VND"));

        ProductInfo result = resilientClient.getProduct(1L);
        assertEquals("Test", result.name());

        // Verify 3 calls made
        verify(productClient, times(3)).getById(1L);
    }

    @Test
    void testFallbackWhenAllRetriesFail() {
        when(productClient.getById(1L))
            .thenThrow(new IOException("down"));

        // Should return fallback (cached or throw)
        // ...
    }
}

演習

1.追加 @Retry + @Timeout 製品サービス REST クライアント呼び出し用 2.実装する @CircuitBreaker フォールバックを使用するとキャッシュされたデータが返されます 3.追加 @Bulkhead 製品サービスへの同時呼び出しを制限する 4.作成 ResilientProductClient すべてのパターンのラッパー 5.作成 ResilientPaymentClient Kafka へのフォールバック キューを使用する 6. 設定を上書きします。 application.properties 7. Micrometer メトリクスを追加します。メトリクスを次の場所で確認します。 /q/metrics 8. テスト: 製品サービスのモックダウン → 再試行の検証 → 回線オープン → フォールバック 9. (上級) フォールト トレランス メトリクス用の Grafana ダッシュボードを作成する

概要

  • @Retry - 自動再試行、指数バックオフのサポート
  • @Timeout — 処理時間を制限し、スレッドのブロックを防ぎます
  • @CircuitBreaker — サービスが失敗しすぎるとサーキット ブレーク (クローズ → オープン → ハーフオープン)
  • @Bulkhead — 同時呼び出しを制限し、リソースの枯渇を回避します
  • @Fallback — すべてが失敗した場合、または再試行をキューに置いた場合に、キャッシュされた/デフォルトのデータを返します。
  • @RateLimit — リクエスト レートの制限 (リクエスト/ウィンドウ)
  • 実行順序: バルクヘッド → CircuitBreaker → 再試行 → タイムアウト → メソッド → フォールバック
  • カスケード障害防止 — タイムアウト + サーキットブレーカー + バルクヘッドを組み合わせます
  • メトリクス — Micrometer/Prometheus 経由で自動的に公開されます
  • 構成オーバーライド が渡されました application.properties — コードを変更する必要はありません

次の記事: OpenTelemetry — 分散トレーシングとメトリクス。