30ステージのパイプラインを「どこで落ちても再開できる」設計にした話
はじめに
業務自動化で「20以上のAPIを順番に叩く」パイプラインを組んだことがある方は多いと思います。最初は await を並べるだけで動きますが、本番で3ヶ月運用すると確実にこう聞かれます。
「途中で落ちたんだけど、最初からやり直し?」
私が構築したローカルで動く業務自動化システムでは、約30ステージのパイプラインを設計しました。本記事では、このパイプラインを「どのステージで落ちても安全に再開できる」設計にするために実装したパターンを紹介します。
前提・環境
- TypeScript + Node.js
- PocketBase(SQLite ベースの軽量DB)を状態管理に使用
- 各ステージは外部APIやKubernetesクラスターへのリモートコマンドを実行
パイプラインの全体像

まず、メインの run() メソッドを見てください。
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);
}
}
このパターンの重要なポイントは以下です。
ステータスでフィルタリングしてから処理する
各ステージは「自分の前のステータス」を持つレコードだけを処理します。例えば stage2 は PROFILE_ATTACHED のレコードだけを取得します。これにより、パイプラインを何度再実行しても、すでに完了したステージはスキップされます。
リトライカウンターの3段階制御

| メソッド | 用途 |
|---|---|
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 を並べるだけ」のパイプラインから、「どこで落ちても安全に再開できる」パイプラインにするために必要な要素をまとめます。
- ステージごとの状態永続化 — 各ステージの完了をDBに記録し、ステータスでフィルタリングして処理する
- 二重実行防止 — パイプラインレベルとステージレベルの2段階ガード
- リトライ制御 — 一時的失敗はリトライ、永続的失敗は停止
- タイプ別分岐 —
run()は全ステージを直列呼び出し、各ステージ内でフィルタ - 監査ログ — 全入出力の記録
- エラー通知 — 失敗を即座にチームへ通知
特別なフレームワークは不要です。PocketBase(またはSQLite)とシンプルなステータス管理だけで、十分に堅牢なパイプラインを構築できます。








