0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

TypeScriptでBulkheadを実装する:依存先ごとにPromise Poolを分離して連鎖遅延を防ぐ

0
Posted at

先に結論

外部依存を呼ぶ非同期処理では、依存先ごとに並行数を分けると、一方の遅延が他方の待ち行列を埋めることを防げます。これが Bulkhead(隔壁)パターンです。

この記事では、Node.js / TypeScript だけで次を満たす小さな実装を作ります。

  • STT、LLM、検索 API のような依存先ごとに上限を持てる
  • 待機中のリクエストを AbortSignal で取り消せる
  • ある依存先が詰まっても、別の依存先の枠は使える
  • テストで「隔離されている」ことを確認できる

リアルタイム処理では、全 API に一つの Promise Pool を共有すると、遅い依存先がワーカーを占有し、関係のない処理まで待たせます。並行数の合計を減らすだけでは、この連鎖を止められません。

これはレートリミットではありません。レートリミットは「単位時間あたりの開始数」を制御し、Bulkhead は「同時に実行してよい数」と障害範囲を制御します。両方が必要なケースもあります。

まず決めるべき境界

Bulkhead のキーは、URL 全体ではなく同じ障害モードを共有する依存先にします。たとえば同じ LLM ベンダーでも、短い分類 API と長い生成 API でタイムアウトやコストの性質が違うなら、別の Pool にします。

依存先 上限の例 分ける理由
STT 2 音声長により実行時間が伸びやすい
LLM 生成 1 ストリームが長く、同時実行のコストが高い
検索 API 3 短い I/O が中心で、別障害として扱える

上限の値は「サーバーの CPU コア数」から機械的に決めません。相手 API の同時接続制限、p95 レイテンシ、リトライ時の増幅、ユーザーごとの公平性を見て調整します。

可取消な Bulkhead を実装する

キューに入った時点では枠を消費しません。枠を得た直後に running を増やし、タスクが成功・失敗・中断のどれで終わっても finally で必ず返します。

type Task<T> = (signal: AbortSignal) => Promise<T>;

const abortReason = (signal: AbortSignal) =>
  signal.reason ?? new DOMException("Aborted", "AbortError");

class Bulkhead {
  private running = 0;
  private readonly waiting: Array<() => void> = [];

  constructor(private readonly limit: number) {
    if (!Number.isInteger(limit) || limit < 1) {
      throw new RangeError("limit must be a positive integer");
    }
  }

  async run<T>(
    task: Task<T>,
    signal = new AbortController().signal,
  ): Promise<T> {
    await this.acquire(signal);

    try {
      if (signal.aborted) throw abortReason(signal);
      return await task(signal);
    } finally {
      this.release();
    }
  }

  private acquire(signal: AbortSignal): Promise<void> {
    if (signal.aborted) return Promise.reject(abortReason(signal));

    if (this.running < this.limit) {
      this.running += 1;
      return Promise.resolve();
    }

    return new Promise((resolve, reject) => {
      const onAbort = () => {
        const index = this.waiting.indexOf(grant);
        if (index !== -1) this.waiting.splice(index, 1);
        reject(abortReason(signal));
      };

      const grant = () => {
        signal.removeEventListener("abort", onAbort);
        this.running += 1;
        resolve();
      };

      signal.addEventListener("abort", onAbort, { once: true });
      this.waiting.push(grant);

      // addEventListener 後に abort された場合も取りこぼさない。
      if (signal.aborted) onAbort();
    });
  }

  private release() {
    this.running -= 1;
    this.waiting.shift()?.();
  }
}

この実装の重要な不変条件は二つです。

  1. 0 <= running <= limit を常に守ること
  2. 待機中に中断されたタスクはキューから取り除き、実行枠を取らないこと

finally を省くと、例外を投げた一件が枠を返さず、Pool は時間とともに停止します。実運用で最初に確認したい箇所です。

依存先ごとに Registry を置く

呼び出し側が Pool を直接選ぶより、依存先名をキーにする Registry にすると、上限の一覧を一か所に置けます。

class BulkheadRegistry {
  private readonly pools = new Map<string, Bulkhead>();

  constructor(limits: Record<string, number>) {
    for (const [name, limit] of Object.entries(limits)) {
      this.pools.set(name, new Bulkhead(limit));
    }
  }

  run<T>(name: string, task: Task<T>, signal?: AbortSignal) {
    const pool = this.pools.get(name);
    if (!pool) throw new Error(`unknown dependency: ${name}`);
    return pool.run(task, signal);
  }
}

const dependencies = new BulkheadRegistry({
  stt: 2,
  llm: 1,
  search: 3,
});

const controller = new AbortController();

const answer = await dependencies.run(
  "llm",
  async (signal) => {
    if (signal.aborted) throw abortReason(signal);
    return { text: "generated" };
  },
  controller.signal,
);

AbortSignal は Pool の待機だけでなく、実際の HTTP 呼び出しでは fetch にも渡します。待機中に期限切れならキューから外れ、すでに実行中なら下流 I/O も中断できます。片方だけでは、不要になった仕事が残ります。

隔離をテストする

「STT の待ち行列があるのに LLM が完了する」ことを、時間依存のテストではなくゲートで確認します。以下は Node.js の assert と Bun / Node.js のどちらでも実行できます。

import assert from "node:assert/strict";

const registry = new BulkheadRegistry({ stt: 1, llm: 1 });

let releaseStt!: () => void;
const sttGate = new Promise<void>((resolve) => {
  releaseStt = resolve;
});
let queuedSttStarted = false;

const firstStt = registry.run("stt", async () => {
  await sttGate;
  return "first";
});
await Promise.resolve();

const queuedStt = registry.run("stt", async () => {
  queuedSttStarted = true;
  return "second";
});

// STT は満席でも、別 Pool の LLM は待たない。
const llmResult = await registry.run("llm", async () => "LLM is isolated");
assert.equal(llmResult, "LLM is isolated");
assert.equal(queuedSttStarted, false);

// 待機中の STT は中断後に実行されない。
const aborter = new AbortController();
const canceledStt = registry.run(
  "stt",
  async () => "must not run",
  aborter.signal,
);
aborter.abort(new Error("request deadline exceeded"));
await assert.rejects(canceledStt, /deadline exceeded/);

releaseStt();
assert.equal(await firstStt, "first");
assert.equal(await queuedStt, "second");
assert.equal(queuedSttStarted, true);

このテストで検証しているのは処理時間ではなく、障害境界です。STT の一件目を意図的に止めても LLM の完了を待たせないため、CI の負荷による揺れが入りません。

運用で追加するもの

この最小実装をサービスに入れるなら、少なくとも次の観測値を追加します。

  • 依存先別の running と待機数
  • キュー待機時間の p95 / p99
  • 中断数と中断理由
  • 下流エラー数・タイムアウト数
  • 上限に達して拒否した件数

待機数だけを見て上限を増やすのは危険です。下流が遅いのに上限を上げると、接続・メモリ・リトライが同時に増え、障害を広げます。まず「どの依存先が詰まり、どれだけ待ち、いつ中断されたか」を分けて観測します。

面接で説明するなら

設計面接では、次の順で説明すると意図が伝わりやすくなります。

  1. 外部依存ごとに障害モードが異なるため、実行枠を共有しない
  2. タスクは FIFO で待たせ、待機中のキャンセルはキューから除去する
  3. finally で必ず枠を返し、成功・失敗・中断で状態を揃える
  4. Pool ごとの待機時間と拒否数を観測し、下流の容量に合わせて上限を調整する

Bulkhead は高速化のための小技ではなく、遅い依存先を「遅いまま閉じ込める」ための境界です。共有 Pool を一つ増やす前に、どの依存先が同じ失敗をするのかを見直すと、障害時の挙動がかなり読みやすくなります。

0
0
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?