堅牢なパイプライン — 例外、再試行、サーキットブレーカーにより、500回の呼び出しを安全に処理
このトピックを終えると
教科書で学んだ 例外処理、try-except、エラー処理 を組み合わせて、外部 API を大量に呼び出すパイプラインが、いくつかの失敗で停止することなく、最後まで実行できるようにするツールを自分で作成できます。成功率 95% を 99.9% に向上させる実戦的なパターンです。
この記事は 教育目的の一般的な例 です。実戦では、tenacity、backoff、resilience4j などのライブラリを使用します。
「次の日の朝、ログを開いた」— 素朴なパイプラインの落とし穴
あなたが500個の遺伝子について、NCBI BLASTを順番に呼び出すスクリプトを夜通し実行したとしましょう。
def naive_pipeline(gene_ids: list[str]) -> list[dict]: results = [] for gene_id in gene_ids: response = call_blast_api(gene_id) results.append(response) return results次の日の朝、ログを開くと、次のようになっています。
[00:03] gene 1: OK
[00:05] gene 2: OK
...
[01:23] gene 27: TIMEOUT
Traceback (most recent call last):
...
requests.exceptions.Timeout全体のスクリプトは、27番目の呼び出しで中断され、残りの473回は試行されませんでした。1つの失敗が全体の実行を停止させたのです。
このアプローチの問題点:
問題1: 失敗の分離がない。1回のエラーが全体を停止させる。
問題2: 再試行がない。ネットワークの障害は一時的なことが多い。少し待ってから再試行すれば、成功する可能性が高くなります。
問題3: 進行状況の保存がない。成功した26個の結果がメモリに保存されたまま、クラッシュ。最初からやり直し。
適切なアプローチは、堅牢なパイプラインパターンです。各呼び出しの失敗を分離し、再試行を行い、永続的な失敗の場合は処理を停止します。成功した結果は、すぐに保存します。
ブラックボックスからコンポーネントへ
コンポーネント 1: try-except で失敗を分離
def robust_pipeline_basic(gene_ids: list[str]) -> tuple[list[dict], list[dict]]: results = [] errors = [] for gene_id in gene_ids: try: response = call_blast_api(gene_id) results.append({"gene_id": gene_id, "data": response}) except Exception as e: errors.append({"gene_id": gene_id, "error": type(e).__name__, "message": str(e)}) return results, errorsポイント: 失敗した呼び出しは errors に分離し、次の反復処理を続行します。 500 個中、いくつかの処理が失敗しても、残りは完了します。
注意: except Exception はすべての例外を捕捉します。 これは、パイプラインの最上位でのみ使用してください。 特定のエラーのみを捕捉したい場合は、明示的に指定します。
try: response = call_blast_api(gene_id)except requests.exceptions.Timeout: # 再試行可能なエラー passexcept requests.exceptions.HTTPError as e: if e.response.status_code == 429: # レート制限 → 再試行可能 pass else: # その他の 4xx → 再試行しても意味がない raiseコンポーネント 2: 指数バックオフによる再試行
一時的なエラーは、再試行によってほとんど解決できます。 ただし、すぐに再試行するとサーバーに負荷がかかります。 したがって、各再試行の間隔を指数関数的に増加させます。
import timeimport random
def with_retry(func, max_attempts: int = 5, base_delay: float = 1.0): """ func を呼び出し、失敗した場合は指数バックオフで再試行します。 """ for attempt in range(1, max_attempts + 1): try: return func() except (requests.exceptions.Timeout, requests.exceptions.ConnectionError) as e: if attempt == max_attempts: raise delay = base_delay * (2 ** (attempt - 1)) jitter = random.uniform(0, delay * 0.1) time.sleep(delay + jitter)間隔: 1秒、2秒、4秒、8秒、16秒。 5回の再試行で、最大31秒間待機します。
ジッター(ランダムノイズ) を加える理由: 複数のクライアントが同時に失敗し、正確に同じタイミングで再試行すると、サーバーが再びダウンする可能性があります。 各クライアントがわずかに異なるタイミングで再試行するように、ランダムノイズを加えます。
コンポーネント 3: サーキットブレーカー
サーバーが長期間ダウンしている場合、再試行しても意味がありません。 連続した失敗が閾値を超えた場合、一時的に呼び出しを完全に停止します。
from dataclasses import dataclassfrom enum import Enum
class CircuitState(Enum): CLOSED = "closed" # 正常 - 呼び出しを許可 OPEN = "open" # 遮断 - すぐに失敗 HALF_OPEN = "half_open" # 回復の確認 - 試験的な呼び出し
@dataclassclass CircuitBreaker: failure_threshold: int = 5 recovery_timeout: float = 60.0 state: CircuitState = CircuitState.CLOSED failure_count: int = 0 last_failure_time: float = 0.0 def call(self, func): now = time.time() if self.state == CircuitState.OPEN: if now - self.last_failure_time > self.recovery_timeout: self.state = CircuitState.HALF_OPEN else: raise Exception("Circuit breaker OPEN") try: result = func() if self.state == CircuitState.HALF_OPEN: self.state = CircuitState.CLOSED self.failure_count = 0 return result except Exception: self.failure_count += 1 self.last_failure_time = now if self.failure_count >= self.failure_threshold: self.state = CircuitState.OPEN raise3つの状態:
- CLOSED: 正常。呼び出しを許可します。
- OPEN: 遮断。すぐに失敗を返します(実際の呼び出しは行いません)。
- HALF_OPEN: 回復を試みます。1つの呼び出しでサーバーの状態をテストします。
パイプラインの組み立て
3つの部品を組み合わせます。
def robust_pipeline_full( gene_ids: list[str], output_path: str, checkpoint_every: int = 50) -> dict: breaker = CircuitBreaker(failure_threshold=5, recovery_timeout=60.0) results = [] errors = [] # 既存のチェックポイントをロード already_done = load_checkpoint(output_path) for i, gene_id in enumerate(gene_ids): if gene_id in already_done: continue try: def call(): return breaker.call(lambda: call_blast_api(gene_id)) response = with_retry(call, max_attempts=3, base_delay=2.0) results.append({"gene_id": gene_id, "data": response}) except Exception as e: errors.append({ "gene_id": gene_id, "error": type(e).__name__, "message": str(e) }) if (i + 1) % checkpoint_every == 0: save_checkpoint(output_path, results) print(f"Progress: {i+1}/{len(gene_ids)} — checkpoint saved") save_checkpoint(output_path, results) return { "total": len(gene_ids), "success": len(results), "failed": len(errors), "errors": errors }
def save_checkpoint(path: str, results: list) -> None: import json with open(path, "w") as f: json.dump(results, f, indent=2)
def load_checkpoint(path: str) -> set[str]: import json from pathlib import Path if not Path(path).exists(): return set() with open(path) as f: data = json.load(f) return {r["gene_id"] for r in data}これで、クラッシュが発生した場合でも、再度実行すると、すでに処理済みのgeneはスキップされ、残りの処理のみが行われます。
フェーディング — 埋めるべき3つの空白
空白1:例外ごとの再試行ポリシー
一部の例外のみ再試行する必要があります。たとえば、400 Bad Requestは再試行しても意味がありません。
def with_retry_smart(func, max_attempts: int = 5): retryable = ( requests.exceptions.Timeout, requests.exceptions.ConnectionError, ) for attempt in range(1, max_attempts + 1): try: return func() except retryable as e: # TODO: 再試行ロジック(バックオフ+ジッター) pass except requests.exceptions.HTTPError as e: # TODO: status_codeを確認し、429または5xxの場合は再試行 # それ以外の場合はすぐにraise passヒント: if e.response.status_code == 429 or 500 <= e.response.status_code < 600: sleep(delay); continue; else: raise.
空白2:失敗レポート
失敗したgene_idsを分析し、原因ごとに分類します。
def analyze_failures(errors: list[dict]) -> dict: """ エラーリストを種類ごとにグループ化します。 返却値:{"Timeout": 5, "HTTPError_500": 3, ...} """ from collections import Counter # TODO: 各エラーの"error"フィールドでカウント # HTTPErrorはstatus_codeで細分化 passヒント: counter = Counter(); for e in errors: key = e["error"]; if "status" in e["message"]: key = f"{key}_{parse_status}"; counter[key] += 1.
空白3:再実行時に失敗した部分のみ再試行
チェックポイントに成功/失敗の両方を保存し、再実行時に失敗した部分のみ再試行します。
def retry_failed_only( gene_ids: list[str], output_path: str, error_log_path: str) -> None: # TODO 1: error_log_pathをロード # TODO 2: 失敗したgene_idsのみを抽出 # TODO 3: robust_pipeline_fullを再実行(フィルタリングされたリストを使用) pass考察:実戦における堅牢性ツールの違い
tenacityライブラリ: Pythonの標準的な再試行ツール。デコレータにより、より簡潔な構文が可能。
from tenacity import retry, stop_after_attempt, wait_exponential
@retry(stop=stop_after_attempt(5), wait=wait_exponential(multiplier=1, max=60))def call_api(gene_id): return requests.get(f"...").json()Kubernetes readiness probes: 実戦環境でのデプロイでは、アプリケーション自体が自身の状態を監視する。問題が検出された場合、トラフィックを別の場所にリダイレクトする。
Netflix Hystrix / resilience4j: Javaエコシステムの洗練されたサーキットブレーカー。バルクヘッド、タイムリミッター、レートリミッターなどのパターンを組み合わせる。
Dead Letter Queue: 失敗したリクエストを別のキューに保存し、後で手動または自動で再処理する。AWS SQS、RabbitMQの標準機能。
分散トランザクションの代替手段: 実戦では、失敗しても回復可能な方法(Sagaパターン、イベントソーシング)で設計する。失敗を例外ではなく、通常のケースとして扱う。
拡張プロジェクト
1. tenacityによる再実装: 既存のバックオフロジックをtenacityデコレータに置き換える。
2. 並列処理: concurrent.futures.ThreadPoolExecutor を使用して、複数の呼び出しを同時に実行する。サーキットブレーカーは複数のスレッド間で共有される。
3. ダッシュボード: 進行状況、成功率、サーキットブレーカーの状態をリアルタイムで表示する。Streamlitを使用すれば30分で作成可能。
4. 通知: 失敗率が閾値を超えた場合に、Slack/Telegramで通知する。
このセクションの構成要素
- [F] 例外: Pythonの階層的な例外システム。
Exceptionを継承し、特定のタイプのみを捕捉します。 - [F] try-except: 失敗の隔離と回復。
except ... as eでエラー情報にアクセスします。 - [F] エラー処理: 再試行、バックオフ、サーキットブレーカーを組み合わせたパターン。
- [W] ファイルI/O: チェックポイントの保存と読み込み(完全なスクリプトを提供)。
[F] = 自分で実装するもの / [W] = 完全なコードで提供するもの。