概要
Circuit Breaker は「電気ブレーカー」と同じ発想
連続失敗を検知したら回路を OPEN にして遮断し、障害源(ベンダーAPI)への負荷と、呼び出し元(Cloud Run の同時実行枠)の両方を保護する。CLOSED → OPEN → HALF_OPEN → CLOSED の状態遷移を StrEnum + match 文で明示的にモデリングする。
タイムアウトはリトライより先に効く
Bad実装は httpx.post にタイムアウトがなく無限待機する。asyncio.wait_for で強制タイムアウト(3秒)を掛けることで、リトライ戦略が機能する前提を作る。フェイルファストの第一歩。
エラーバジェット = 許容される失敗の予算
SLO 99.9%・30日ローリングウィンドウなら許容ダウンタイムは 43.2分。「今どれだけ予算を消費したか」を定量化すると、開発速度とのバランスを数値で語れるようになる。
バーンレート = 予算消費の「速度」を見る
固定閾値(5分窓・1%)は「今エラーが多いか」しか見ない。バーンレートは「このペースが続くとSLOをいつ違反するか」を見る。fast-burn(page)と slow-burn(ticket)で重大度を分離し、アラート疲れを防ぐ。
問題
ECサイト MOps チームの「キャンペーン配信 API」(Cloud Run v2)は、SMS/プッシュ通知を送るために外部ベンダーの通知 API を同期呼び出ししている。先週、ベンダー側で障害(5xx 急増・レスポンス遅延)が発生した際、以下の「悪いコード」がベンダーへのリクエストを無制限にリトライし続けたため、Cloud Run の同時実行枠を使い切り、ベンダーとは無関係な健全なキャンペーンの配信リクエストまで巻き込まれてタイムアウトするというカスケード障害が発生した。さらに、監視アラートは固定閾値のみで運用されており、障害検知までに 40 分かかった。
制約・前提条件
- Python 3.12+、
asyncio+httpx.AsyncClient - ベンダー API SLA: 99.5%(月間)だが、障害時はさらに悪化する
- 自社 SLO 目標: キャンペーン配信 API の可用性 99.9%(30日ローリングウィンドウ)
- Cloud Run v2 の同時実行数(
--concurrency)は有限で、ブロッキング待機はこの枠を消費する - Circuit Breaker パラメータ:
failure_threshold=5/recovery_timeout_sec=30/half_open_max_calls=3/half_open_success_threshold=2 - リクエストタイムアウト: 3秒(現状は未設定 → 無限待機)
- アラートは DataDog Monitor を想定し、Google SRE Workbook のマルチウィンドウ・マルチバーンレートアラート方式(fast-burn / slow-burn の2段階)を採用する
悪いコード (Before)
import time
import httpx
VENDOR_API_URL = "https://vendor-sms.example.com/v1/send"
def send_sms(phone: str, message: str) -> bool:
# 問題①: タイムアウト未設定 → ベンダー障害時に無限待機
while True:
try:
resp = httpx.post(VENDOR_API_URL, json={"to": phone, "body": message})
if resp.status_code == 200:
return True
# 問題②: 5xxでも同じロジックでリトライ継続(回数無制限・固定1秒)
time.sleep(1)
except httpx.RequestError:
time.sleep(1)
continue
# 問題③: Circuit Breakerなし → 障害中も全リクエストが送られ続ける
# 問題④: 呼び出し側(Cloud Run)を保護する仕組みがない
# → 同時実行枠を使い切りカスケード障害
def check_alert(error_count: int, total_count: int) -> None:
# 問題⑤: 固定閾値のみ(5分窓・1%)
# → 小さいスパイクでもページが鳴りアラート疲れ
error_rate = error_count / total_count if total_count else 0
if error_rate > 0.01:
page_oncall()
# 問題⑥: 同じ固定閾値では6時間かけて進行する
# 低速バーンを検知できない → 検知が遅れる
# 問題⑦: エラーバジェットという定量指標が存在せず、
# 開発速度とのバランスを判断できない
def page_oncall() -> None:
...
asyncio.wait_for で強制ヒント(段階的開示)
ヒント1 — 方向性
ヒント2 — アプローチ
- 問題①②③④:
CircuitStateをStrEnumでCLOSED/OPEN/HALF_OPENと定義。CLOSEDで連続failure_threshold回失敗→OPEN。OPENでrecovery_timeout_sec経過→HALF_OPEN。HALF_OPENで連続成功→CLOSED、1回でも失敗すれば即OPENに戻す asyncio.wait_for(func(...), timeout=REQUEST_TIMEOUT_SEC)でタイムアウトを強制- リトライは
min(BASE * 2**attempt, MAX) + jitterの指数バックオフに置き換える。ただしOPEN中はリトライせず即座に失敗(フェイルファスト) - 問題⑤⑥⑦: SLO 99.9% × 30日 → エラーバジェット(分) =
30*24*60*(1-0.999)= 43.2分 - マルチバーンレート公式:
burn_rate = budget_consumed_fraction * (720時間 / window_hours)。「1時間で2%消費」なら0.02*720/1=14.4。許容エラー率閾値 =burn_rate * (1 - slo_target)
ヒント3 — コードの骨格
class CircuitState(StrEnum):
CLOSED = "closed"; OPEN = "open"; HALF_OPEN = "half_open"
class CircuitBreaker:
async def call(self, func, *args, **kwargs):
async with self._lock:
self._transition_if_needed() # OPEN→HALF_OPEN 時間経過チェック
if self._state is CircuitState.OPEN:
raise VendorUnavailableError(...)
try:
result = await asyncio.wait_for(func(*args, **kwargs), timeout=REQUEST_TIMEOUT_SEC)
except Exception:
await self._on_failure() # 状態遷移: →OPEN
raise
else:
await self._on_success() # 状態遷移: HALF_OPEN→CLOSED
return result
# SLO/エラーバジェット
SLO_TARGET = 0.999
ERROR_BUDGET_MINUTES = 30 * 24 * 60 * (1 - SLO_TARGET) # 43.2
def _burn_rate(budget_consumed_fraction: float, window_hours: float) -> float:
return budget_consumed_fraction * (30 * 24) / window_hours
問題点分析(7点)
| # | 問題点 | 分類 | 改善方法 |
|---|---|---|---|
| 1 | タイムアウト未設定(無限待機) | 耐障害性 | asyncio.wait_for(timeout=3.0) で一元管理 |
| 2 | 無制限リトライ・固定1秒スリープ | 耐障害性 | 指数バックオフ + jitter(min(BASE*2**n, MAX)) |
| 3 | Circuit Breaker がない | 耐障害性 | CLOSED/OPEN/HALF_OPEN の状態遷移で遮断 |
| 4 | 呼び出し元(Cloud Run)保護機構なし | 可用性 | OPEN 中はフェイルファストで即座に拒否 |
| 5 | 固定閾値アラート(5分窓・1%) | アラート設計 | マルチウィンドウ・マルチバーンレートに置換 |
| 6 | 低速バーンを検知できない | アラート設計 | slow-burn(6h窓・6倍速)を ticket として監視 |
| 7 | エラーバジェット概念の欠如 | 可観測性 | SLO 99.9%→予算43.2分を明示的に定数化 |
Circuit Breaker 状態遷移図 + マルチバーンレート早見表
模範解答
def send_sms(phone: str, message: str) -> bool:
while True:
try:
resp = httpx.post(VENDOR_API_URL, json={"to": phone, "body": message})
if resp.status_code == 200:
return True
time.sleep(1)
except httpx.RequestError:
time.sleep(1)
continue
from __future__ import annotations
import asyncio
import random
import time
from dataclasses import dataclass
from enum import StrEnum
from typing import Callable, Final
import httpx
# ── 定数 ──────────────────────────────────────────
VENDOR_API_URL: Final[str] = "https://vendor-sms.example.com/v1/send"
REQUEST_TIMEOUT_SEC: Final[float] = 3.0
MAX_RETRY_ATTEMPTS: Final[int] = 3
BASE_BACKOFF_SEC: Final[float] = 0.5
MAX_BACKOFF_SEC: Final[float] = 4.0
class CircuitState(StrEnum):
"""Circuit Breaker の状態(3値のみを許容する値オブジェクト)."""
CLOSED = "closed" # 正常時: リクエストを全て通す
OPEN = "open" # 遮断時: リクエストを即座に拒否
HALF_OPEN = "half_open" # 復旧確認中: 限定的にプローブのみ通す
@dataclass(frozen=True, slots=True)
class CircuitBreakerConfig:
"""Circuit Breaker の設定値(完全コンストラクタ・不変)."""
failure_threshold: int = 5
recovery_timeout_sec: float = 30.0
half_open_max_calls: int = 3
half_open_success_threshold: int = 2
class VendorUnavailableError(RuntimeError):
"""Circuit が OPEN 中にリクエストが拒否されたことを示す例外."""
class CircuitBreaker:
"""ベンダー API 呼び出しを保護する Circuit Breaker(asyncio 対応)."""
def __init__(self, config: CircuitBreakerConfig | None = None) -> None:
self._config = config or CircuitBreakerConfig()
self._state = CircuitState.CLOSED
self._failure_count = 0
self._half_open_calls = 0
self._half_open_successes = 0
self._opened_at: float = 0.0
self._lock = asyncio.Lock()
async def call(self, func: Callable, *args: object, **kwargs: object) -> object:
"""保護対象の非同期関数を Circuit Breaker 経由で呼び出す.
Args:
func: 呼び出す非同期関数。
*args, **kwargs: func に渡す引数。
Returns:
func の戻り値。
Raises:
VendorUnavailableError: OPEN で拒否した場合。
"""
async with self._lock:
self._transition_if_needed()
if self._state is CircuitState.OPEN:
raise VendorUnavailableError("circuit is OPEN: vendor API を保護中")
if self._state is CircuitState.HALF_OPEN:
if self._half_open_calls >= self._config.half_open_max_calls:
raise VendorUnavailableError("circuit is HALF_OPEN: プローブ上限")
self._half_open_calls += 1
try:
# タイムアウトを一元管理(Bad実装は無限待機)
result = await asyncio.wait_for(
func(*args, **kwargs), timeout=REQUEST_TIMEOUT_SEC
)
except Exception:
await self._on_failure()
raise
else:
await self._on_success()
return result
def _transition_if_needed(self) -> None:
"""OPEN → HALF_OPEN への時間経過遷移(lock内で呼ぶ前提)."""
match self._state:
case CircuitState.OPEN:
elapsed = time.monotonic() - self._opened_at
if elapsed >= self._config.recovery_timeout_sec:
self._state = CircuitState.HALF_OPEN
self._half_open_calls = 0
self._half_open_successes = 0
case _:
pass
async def _on_success(self) -> None:
async with self._lock:
match self._state:
case CircuitState.HALF_OPEN:
self._half_open_successes += 1
if self._half_open_successes >= self._config.half_open_success_threshold:
self._state = CircuitState.CLOSED
self._failure_count = 0
case CircuitState.CLOSED:
self._failure_count = 0
case _:
pass
async def _on_failure(self) -> None:
async with self._lock:
match self._state:
case CircuitState.HALF_OPEN:
# プローブ失敗 → 即座に OPEN に戻す
self._state = CircuitState.OPEN
self._opened_at = time.monotonic()
case CircuitState.CLOSED:
self._failure_count += 1
if self._failure_count >= self._config.failure_threshold:
self._state = CircuitState.OPEN
self._opened_at = time.monotonic()
case _:
pass
async def _post_vendor_sms(phone: str, message: str) -> bool:
"""ベンダー API への実際の POST(タイムアウトは呼び出し元で一元管理)."""
async with httpx.AsyncClient() as client:
resp = await client.post(VENDOR_API_URL, json={"to": phone, "body": message})
resp.raise_for_status()
return True
_breaker = CircuitBreaker()
async def send_sms(phone: str, message: str) -> bool:
"""SMS を送信する(Circuit Breaker + 指数バックオフ付き).
Args:
phone: 送信先電話番号。
message: 送信本文。
Returns:
送信成功なら True。
Raises:
VendorUnavailableError: OPEN 中で送信不可(リトライしない)。
"""
last_exc: Exception | None = None
for attempt in range(MAX_RETRY_ATTEMPTS):
try:
return await _breaker.call(_post_vendor_sms, phone, message)
except VendorUnavailableError:
raise # OPEN中はリトライせず即座に失敗(フェイルファスト)
except Exception as exc: # noqa: BLE001
last_exc = exc
backoff = min(BASE_BACKOFF_SEC * (2**attempt), MAX_BACKOFF_SEC)
jitter = random.uniform(0, backoff * 0.1)
await asyncio.sleep(backoff + jitter)
raise last_exc # type: ignore[misc]
# slo_alerting.py ── SLO/エラーバジェット × マルチバーンレートアラート設計
"""SLO 99.9%(30日ローリングウィンドウ)に基づくアラート定義.
Google SRE Workbook の「マルチウィンドウ・マルチバーンレートアラート」方式。
固定閾値(Bad実装の「5分窓・1%」)を廃止し、エラーバジェットの
消費速度(バーンレート)に応じて重大度(page/ticket)を分離する。
"""
from __future__ import annotations
from dataclasses import dataclass
from typing import Final
# ── 定数 ──────────────────────────────────────────
SLO_TARGET: Final[float] = 0.999 # 30日ローリングウィンドウで 99.9%
WINDOW_DAYS: Final[int] = 30
WINDOW_HOURS: Final[float] = WINDOW_DAYS * 24 # 720時間
# エラーバジェット = 許容ダウンタイム(分)
ERROR_BUDGET_MINUTES: Final[float] = WINDOW_DAYS * 24 * 60 * (1 - SLO_TARGET) # 43.2分
@dataclass(frozen=True, slots=True)
class BurnRateAlert:
"""マルチウィンドウ・マルチバーンレートアラート定義(不変の値オブジェクト).
Attributes:
name: アラート名。
long_window_hours: バーンレート判定に使う長期ウィンドウ(時間)。
short_window_minutes: 短期確認ウィンドウ(分)。誤検知抑制用。
burn_rate: このウィンドウで許容する最大バーンレート(倍率)。
budget_consumed_fraction: long_window 内で消費されるバジェット割合。
severity: "page"(即時呼び出し)または "ticket"(翌営業日対応)。
"""
name: str
long_window_hours: float
short_window_minutes: float
burn_rate: float
budget_consumed_fraction: float
severity: str
@property
def error_rate_threshold(self) -> float:
"""このバーンレートに対応する許容エラー率閾値を算出する."""
return self.burn_rate * (1 - SLO_TARGET)
def _burn_rate(budget_consumed_fraction: float, window_hours: float) -> float:
"""指定ウィンドウでバジェットを budget_consumed_fraction 消費するバーンレートを計算する.
Args:
budget_consumed_fraction: ウィンドウ内で消費するバジェット割合(例: 0.02 = 2%)。
window_hours: ウィンドウの長さ(時間)。
Returns:
バーンレート(通常のバジェット消費速度に対する倍率)。
"""
return budget_consumed_fraction * WINDOW_HOURS / window_hours
# fast-burn: 1時間でバジェットの2%を消費するペース(約2.9日で枯渇)
_FAST_BURN_RATE: Final[float] = _burn_rate(budget_consumed_fraction=0.02, window_hours=1) # 14.4
# slow-burn: 6時間でバジェットの5%を消費するペース(約5日で枯渇)
_SLOW_BURN_RATE: Final[float] = _burn_rate(budget_consumed_fraction=0.05, window_hours=6) # 6.0
ALERTS: Final[tuple[BurnRateAlert, ...]] = (
BurnRateAlert(
name="fast-burn-page",
long_window_hours=1,
short_window_minutes=5,
burn_rate=_FAST_BURN_RATE,
budget_consumed_fraction=0.02,
severity="page",
),
BurnRateAlert(
name="slow-burn-ticket",
long_window_hours=6,
short_window_minutes=30,
burn_rate=_SLOW_BURN_RATE,
budget_consumed_fraction=0.05,
severity="ticket",
),
)
# 実行例:
# >>> ERROR_BUDGET_MINUTES
# 43.2
# >>> [a.error_rate_threshold for a in ALERTS]
# [0.0144, 0.006] # fast-burn: 1.44%(1h窓)超えでpage / slow-burn: 0.6%(6h窓)超えでticket
| 問題 | 修正内容 | 効果 |
|---|---|---|
| ① タイムアウトなし | asyncio.wait_for(timeout=3.0) | 無限待機を排除しCloud Runの枠を守る |
| ② 無制限リトライ・固定スリープ | 指数バックオフ + jitter(min(BASE*2**n, MAX)) | 障害中のベンダーへの負荷を抑制 |
| ③④ Circuit Breakerなし | CLOSED/OPEN/HALF_OPEN 状態遷移 + フェイルファスト | カスケード障害を構造的に防止 |
| ⑤⑥ 固定閾値アラート | fast-burn(14.4x/1h)・slow-burn(6x/6h)の2段階 | 誤検知抑制と低速劣化の両方を検知 |
| ⑦ エラーバジェット概念なし | SLO 99.9%→予算43.2分を定数化 | 開発速度とのバランスを定量判断 |
ポイント解説
CLOSED(正常)→ 連続失敗で OPEN(遮断)→ recovery_timeout_sec 経過で HALF_OPEN(プローブ)→ 成功が続けば CLOSED に戻り、1回でも失敗すれば即 OPEN に戻る。match 文で状態ごとの遷移ロジックを明示することで if/elif の分岐地獄を回避(良いコード・悪いコードで学ぶ設計入門 Ch8: 条件分岐)。asyncio.Lock による状態遷移の排他制御 — 複数リクエストが同時に失敗を検知して二重に OPEN 遷移を試みる競合を防ぐ。CircuitBreakerConfig は frozen=True, slots=True の値オブジェクトとして完全コンストラクタ化し、実行中に閾値が書き換わらないことを保証(Ch3: カプセル化・完全コンストラクタ)。_post_vendor_sms 自体にはタイムアウトを設定せず、CircuitBreaker.call() 内の asyncio.wait_for に集約することで、Circuit Breaker を経由しないバイパスが生まれない設計にする。VendorUnavailableError はリトライせず即座に呼び出し元へ伝播させる。OPEN 中にリトライしてしまうと、Circuit Breaker を導入した意味(ベンダー保護・同時実行枠保護)が失われる。burn_rate = budget_consumed_fraction × (720時間 / window_hours)。fast-burn(1時間で2%消費=14.4倍速)は約2.9日でバジェットが尽きるペースなので即ページ、slow-burn(6時間で5%消費=6倍速)は約5日で尽きるペースなのでチケット化に留める。実務への応用
MOps キャンペーン配信 API では、ベンダー通知 API・決済 API・在庫連携 API など外部依存を持つすべての同期呼び出しに CircuitBreaker を共通コンポーネント化して適用できる。障害時に自社サービス全体を守る「防波堤」になる。
DataDog Monitor 設計: ALERTS の error_rate_threshold と long_window_hours / short_window_minutes をそのまま DataDog Monitor のクエリ(avg(last_1h):... > 0.0144 かつ avg(last_5m):... > 0.0144)にマッピングし、severity="page" は PagerDuty 即時通知、severity="ticket" は Jira 起票に振り分ける。
エラーバジェットの週次レビュー: 「今週バジェットを何%消費したか」を Looker Studio ダッシュボードで可視化し、消費率が高い週は新機能リリースを一時停止して信頼性改善に工数を振る、という PM 判断の材料になる(Error Budget Policy)。
今日のまとめ
次のステップ
- 発展問題:
CircuitBreakerをキャンペーン配信で使う複数のベンダー(SMS/プッシュ/メール)ごとに独立したインスタンスとして管理するCircuitBreakerRegistryを設計し、ベンダーごとの障害が他ベンダーに波及しないようにする(Bulkhead パターンとの組み合わせ) - 発展問題:
ALERTSに Google SRE Workbook の標準的な4段階(1h/5m, 6h/30m, 1d/2h, 3d/6h)を追加し、severityを page/ticket の2値から page/ticket/review の3値に拡張する - 参考: Google SRE Workbook「Alerting on SLOs」章、Release It!(Michael T. Nygard)Circuit Breaker パターン、良いコード・悪いコードで学ぶ設計入門 Ch3(カプセル化)/ Ch8(条件分岐)