JavaScript

JavaScript上級:Classで並列数を制御する非同期タスクキューを作る

この記事でわかること

JavaScriptのClassとPromiseで同時実行数を制限するタスクキューを実装。同期例外、非同期の失敗、完了待ち、受付終了、デッドロックの注意点を解説します。

最終回では、同時に実行する非同期処理の数を制限するTaskQueueを作ります。privateフィールドで実行状態を守り、同期例外とPromiseの失敗を同じ経路で処理し、すべてのタスクの完了を待てる設計にします。非同期処理クラス入門を理解してから取り組んでください。

キューの仕様を決める

開始前の関数を受け取る

キューへ渡すのは、実行済みのPromiseではなく、呼び出すと処理を始める関数です。queue.add(() => fetch(url))なら開始時刻をキューで制御できます。先にfetch(url)を呼んでしまうと、登録前に通信が始まります。

守るルール

  • 同時に実行する件数はconstructorへ指定した上限以下です。
  • 待機している関数は登録順に開始します。完了順は保証しません。
  • あるタスクが失敗しても、後続タスクを実行します。
  • close()後は新規登録を拒否し、登録済みの処理は最後まで実行します。
  • onIdle()は、実行中・待機中がともに0になった時点で解決します。

TaskQueueクラスを実装する

task-queue.mjsを作成する

サポート中のNode.js LTSで実行できます。外部ライブラリは使いません。

export class QueueClosedError extends Error {
  constructor() {
    super('キューは受付を終了しています。');
    this.name = 'QueueClosedError';
  }
}

export class TaskQueue {
  #limit;
  #running = 0;
  #waiting = [];
  #idleWaiters = [];
  #closed = false;

  constructor(concurrency = 2) {
    if (!Number.isSafeInteger(concurrency) || concurrency < 1) {
      throw new RangeError('同時実行数は1以上の安全な整数にしてください。');
    }
    this.#limit = concurrency;
  }

  get running() { return this.#running; }
  get pending() { return this.#waiting.length; }

  add(task) {
    if (this.#closed) return Promise.reject(new QueueClosedError());
    if (typeof task !== 'function') {
      return Promise.reject(new TypeError('実行する関数を渡してください。'));
    }
    return new Promise((resolve, reject) => {
      this.#waiting.push({ task, resolve, reject });
      this.#drain();
    });
  }

  close() {
    this.#closed = true;
  }

  onIdle() {
    if (this.#running === 0 && this.#waiting.length === 0) {
      return Promise.resolve();
    }
    return new Promise((resolve) => this.#idleWaiters.push(resolve));
  }

  #drain() {
    while (this.#running < this.#limit && this.#waiting.length > 0) {
      const entry = this.#waiting.shift();
      this.#running += 1;
      Promise.resolve()
        .then(() => entry.task())
        .then(
          (value) => {
            entry.resolve(value);
            this.#finish();
          },
          (error) => {
            entry.reject(error);
            this.#finish();
          },
        );
    }
  }

  #finish() {
    this.#running -= 1;
    this.#drain();
    if (this.#running === 0 && this.#waiting.length === 0) {
      const waiters = this.#idleWaiters;
      this.#idleWaiters = [];
      for (const resolve of waiters) resolve();
    }
  }
}

実行枠を先に確保する

実行を予約した時点で#runningを増やしてから、Promiseのコールバックでタスクを呼びます。先に枠を確保しないと、処理が始まるまでの間に上限以上のタスクを予約してしまうことがあります。

同期例外もPromiseの失敗として扱う

Promise.resolve().then(() => entry.task())を使うと、タスクがその場でthrowした場合も、Promiseを返して後から失敗した場合も、後続の失敗側の処理へ流れます。成功・失敗の両方で#finish()を呼び、実行枠を解放します。

複数タスクを実行する

queue-demo.mjsを作成する

同じフォルダーに保存し、node queue-demo.mjsで実行します。タイマーは時間のかかるI/Oを模したものです。CPU計算を別スレッドで実行しているわけではありません。

import { TaskQueue } from './task-queue.mjs';

const wait = (ms) => new Promise((resolve) => setTimeout(resolve, ms));
const queue = new TaskQueue(2);

function createTask(name, ms, shouldFail = false) {
  return async () => {
    console.log(`開始:${name}`);
    await wait(ms);
    if (shouldFail) throw new Error(`${name}に失敗しました。`);
    console.log(`完了:${name}`);
    return `${name}の結果`;
  };
}

const jobs = [
  queue.add(createTask('A', 30)),
  queue.add(createTask('B', 10, true)),
  queue.add(createTask('C', 5)),
  queue.add(() => { throw new Error('Dの同期エラー'); }),
  queue.add(createTask('E', 5)),
];
queue.close();

const results = await Promise.allSettled(jobs);
await queue.onIdle();
console.log(results.map((result) => result.status));
// fulfilled, rejected, fulfilled, rejected, fulfilled
console.log(queue.running, queue.pending); // 0 0

try {
  await queue.add(createTask('F', 1));
} catch (error) {
  console.log(error.name); // QueueClosedError
}

失敗したPromiseを放置しない

addはタスクごとのPromiseを返します。onIdleはキューが空になったことを通知するだけで、各タスクの成功・失敗を代わりに受け取るものではありません。戻り値をawaitPromise.allSettledで扱い、拒否されたPromiseを放置しないでください。

closeの意味を限定する

close()は新規受付の終了です。待機中のタスクの削除や実行中の処理の中断は行いません。キューが空のときのonIdle()は直ちに解決するため、その後に追加した仕事まで待つわけでもありません。受付終了後に完了を待つと、処理の区切りが明確になります。

難しい処理を検証する

成功件数だけで判断しない

  • 同時実行数が1の場合、開始順と処理の実行順を確認します。
  • 上限が2の場合、実行中カウンターの最大値が2を超えないことを確認します。
  • 同期例外と非同期の失敗の後も、後続タスクが開始されることを確認します。
  • 複数のonIdle()待機がすべて解決することを確認します。
  • 無効な上限や、関数でない引数を拒否することを確認します。

タスク自身がキューの空きを待つ構造を避ける

上限1のキューで実行中のタスクが、同じキューへ別のタスクを追加して完了を待つと、後続が開始できず待ち続ける場合があります。タスクの依存関係を整理し、実行枠を保持したまま同じキューの空きを待たないようにします。

実運用へ拡張する前に

キャンセルとタイムアウト

この実装はタスクのキャンセル機能を持ちません。通信を中断するならAbortSignalをタスクへ渡すなど、タスク側も協調する設計が必要です。Promise.raceで待つのをやめるだけでは、元の通信や計算を停止できません。

待機件数と処理の保存

待機件数に上限がないため、無制限に登録するとメモリーを消費します。大量登録には受付制限が必要です。配列のshiftを使った簡単な実装なので、大規模なキューではデータ構造も見直します。プロセス終了後の再開、複数サーバー間の共有、永続化は対象外です。

継承より役割の分離を選ぶ場面

キューは順番と実行枠を管理し、実際の仕事は関数として受け取ります。仕事の種類が増えても、キューの派生クラスを増やす必要はありません。クラスは状態のルールを守るために使い、処理の差し替えには関数を組み合わせています。

参考資料と連載の振り返り

基礎とのつながり

変数、配列、関数、クロージャー、Promise、クラスが一つの処理でつながりました。仕組みが追いにくい場合はJavaScriptの記事一覧から必要な章へ戻ってください。

参考資料

y.
WRITTEN BY

y_ymo10

SEの部屋で、JavaScript・TypeScript・React.js・Next.jsの開発ノートを公開しています。

ほかの開発ノートを読む →