30ステージのパイプラインを「どこで落ちても再開できる」設計にした話

30ステージのパイプラインを「どこで落ちても再開できる」設計にした話

約30ステージの業務自動化パイプラインを「どこで落ちても再開できる」設計にした実装パターンを紹介。ステータス駆動の冪等性、二重実行防止、3段階リトライ制御をTypeScript + PocketBaseで実現しました。
2026.07.30

はじめに

業務自動化で「20以上のAPIを順番に叩く」パイプラインを組んだことがある方は多いと思います。最初は await を並べるだけで動きますが、本番で3ヶ月運用すると確実にこう聞かれます。

「途中で落ちたんだけど、最初からやり直し?」

私が構築したローカルで動く業務自動化システムでは、約30ステージのパイプラインを設計しました。本記事では、このパイプラインを「どのステージで落ちても安全に再開できる」設計にするために実装したパターンを紹介します。

前提・環境

  • TypeScript + Node.js
  • PocketBase(SQLite ベースの軽量DB)を状態管理に使用
  • 各ステージは外部APIやKubernetesクラスターへのリモートコマンドを実行

パイプラインの全体像

typescript-multi-stage-pipeline-design-pipeline-flow

まず、メインの run() メソッドを見てください。

pipeline.ts
class Pipeline {
  private isRunning = false;

  public async run() {
    if (this.isRunning) {
      Logger.warning("pipeline is already running. Skipping.");
      return;
    }

    this.isRunning = true;

    try {
      await this.stage1_getProfile();       // 共通
      await this.stage2_reserveResource();  // タイプA専用
      await this.stage3_downloadConfig();   // タイプA専用
      await this.stage4_confirmConfig();    // タイプA専用
      await this.stage5_requestShipment();  // 共通
      await this.stage6_confirmShipment();  // 共通
      await this.stage7_sendNotification(); // 共通
      // ... 20ステージ以上続く
      await this.stage30_archive();         // 共通
    } catch (err) {
      console.error(err);
    } finally {
      this.isRunning = false;
    }
  }
}

ここで注目すべきポイントが3つあります。

ポイント1: 二重実行防止のガード

パイプライン全体には isRunning フラグがありますが、各ステージにも独立した isProcessing フラグが存在します。

class Pipeline {
  private isRunning = false;
  private isProcessing = false;

  public async stage1_getProfile(): Promise<void> {
    if (this.isProcessing) {
      Logger.warning("process is already running. Skipping.");
      return;
    }
    this.isProcessing = true;

    try {
      // ... ステージの処理
    } finally {
      this.isProcessing = false;
    }
  }
}

なぜ2つのフラグが必要なのか?isRunning はパイプライン全体の再入を防ぎます。isProcessing は、パイプラインがcronで定期実行される場合に、前回のステージがまだ完了していない状態で次の実行がトリガーされるケースをガードします。

ポイント2: ステージごとの状態永続化

各ステージは「成功したらステータスを更新、失敗したらリトライカウンターを増やす」という共通パターンに従います。

public async stage1_getProfile(): Promise<void> {
  // 1. 対象ステータスのレコードを取得
  const targets = await DB.filterByStatus("INIT");
  if (targets.length === 0) return;

  for (const target of targets) {
    const log = this.createLog("getProfile");

    try {
      // 2. 外部APIを呼ぶ
      const profile = await ExternalAPI.getProfile(target.id);

      // 3. 結果を保存し、ステータスを進める
      await DB.saveProfile(target.id, profile);
      await DB.updateStatus(target.id, "PROFILE_ATTACHED");
      await DB.resetRetry(target.id);

      log.result = "success";
      Logger.success(`retrieved profile for ${target.id}`);
    } catch (err) {
      // 4. 失敗時はリトライカウンターを増やす
      await DB.incRetry(target.id);
      log.result = "failed";
      this.broadcastError(`failed to get profile for ${target.id}`);
    }

    await DB.writeLog(log);
  }
}

このパターンの重要なポイントは以下です。

ステータスでフィルタリングしてから処理する

各ステージは「自分の前のステータス」を持つレコードだけを処理します。例えば stage2PROFILE_ATTACHED のレコードだけを取得します。これにより、パイプラインを何度再実行しても、すでに完了したステージはスキップされます

リトライカウンターの3段階制御

typescript-multi-stage-pipeline-design-retry-control

メソッド 用途
resetRetry 成功時にカウンターを0に戻す
incRetry 一時的な失敗(ネットワークエラー等)でカウンターを+1
setNoRetry 永続的な失敗(データ不整合等)でリトライ自体を停止

incRetry が一定回数を超えたレコードは自動的にスキップされ、setNoRetry が設定されたレコードは手動介入まで処理されません。これにより、1件の不良データがパイプライン全体を止めることを防ぎます。

ポイント3: タイプ別の分岐

実際のシステムでは、処理対象が複数のタイプに分かれることがよくあります。

public async run() {
  // ...
  await this.stage2_reserveResource();  // タイプA専用
  await this.stage3_downloadConfig();   // タイプA専用
  await this.stage5_requestShipment();  // 共通
  await this.stage28_generateQrCode();  // タイプA専用
  await this.stage29_confirmDelivery(); // タイプB専用
  // ...
}

各ステージの内部でタイプフィルタリングを行います。

public async stage2_reserveResource(): Promise<void> {
  // タイプAだけを対象にする
  const typeAItems = await DB.filterByStatus("PROFILE_ATTACHED", "TYPE_A");
  if (typeAItems.length === 0) {
    this.isProcessing = false;
    return; // タイプBはこのステージをスキップ
  }
  // ...
}

run() メソッド側では全ステージを直列に呼びますが、各ステージが「自分に関係ないレコードは0件→即リターン」するため、実質的にタイプ別のフローが実現されます。

ポイント4: 監査ログの自動記録

全ステージの入出力を記録します。

private createLog(stageName: string): StageLog {
  return {
    stage: stageName,
    targetId: "",
    api: "",
    payload: {},
    response: {},
    result: "",
    timestamp: new Date().toISOString(),
  };
}

APIのリクエストとレスポンスをそのまま保存するため、障害時のデバッグが格段に楽になります。「どのステージで」「何を送って」「何が返ってきたか」がすべてDBに残ります。

ポイント5: エラー通知

失敗時にチャットツールへ即座に通知を飛ばします。

private broadcastError(message: string): void {
  const notification = {
    text: `[Pipeline Error] ${message}`,
    timestamp: new Date().toISOString(),
    operator: this.currentOperator,
  };
  NotificationService.send(notification);
}

この設計の良いところと注意点

良いところ

  • 冪等性: 何度実行しても同じ結果。ステータスで制御しているため。
  • 部分再開: 途中で落ちたら、再実行するだけで未処理のレコードから再開。
  • 障害分離: 1件のエラーが他のレコードの処理を止めない。
  • 監査対応: 全ステージの入出力がログとして残る。

注意点

  • ステータス遷移の管理が複雑になる: 30ステージあると、ステータス値が30種類以上。命名規則を統一しないと混乱する。
  • isProcessing フラグはプロセス内のみ有効: マルチプロセスやマルチインスタンス環境では、DB上のロック機構が必要。
  • リトライ上限の適切な設定が重要: 小さすぎると一時的な障害で停止し、大きすぎると不良データが長時間リソースを占有する。

まとめ

await を並べるだけ」のパイプラインから、「どこで落ちても安全に再開できる」パイプラインにするために必要な要素をまとめます。

  1. ステージごとの状態永続化 — 各ステージの完了をDBに記録し、ステータスでフィルタリングして処理する
  2. 二重実行防止 — パイプラインレベルとステージレベルの2段階ガード
  3. リトライ制御 — 一時的失敗はリトライ、永続的失敗は停止
  4. タイプ別分岐run() は全ステージを直列呼び出し、各ステージ内でフィルタ
  5. 監査ログ — 全入出力の記録
  6. エラー通知 — 失敗を即座にチームへ通知

特別なフレームワークは不要です。PocketBase(またはSQLite)とシンプルなステータス管理だけで、十分に堅牢なパイプラインを構築できます。

この記事をシェアする

関連記事