
Bullmq
- 3 installs
- 2 repo stars
- Updated August 3, 2026
- fandhe-ai/agent-reference-skills
Helps with ai & agent building tasks.
About
bullmq is a Claude Code skill for ai & agent building. It helps solo builders move faster with AI-assisted coding.
- bullmq
- AI & Agent Building
- AI-coding skill
Bullmq by the numbers
- 3 all-time installs (skills.sh)
- Ranked #13,677 of 16,546 AI & Agent Building skills by installs in the Skillselion catalog
- Data as of Aug 3, 2026 (Skillselion catalog sync)
npx skills add https://github.com/fandhe-ai/agent-reference-skills --skill bullmqAdd your badge
Show developers this skill is listed on Skillselion. Paste this into your README.
| Installs | 3 |
|---|---|
| repo stars | ★ 2 |
| Last updated | August 3, 2026 |
| Repository | fandhe-ai/agent-reference-skills ↗ |
What it does
Helps with ai & agent building tasks.
Files
BullMQ リファレンス
BullMQ (Redis ベースの Node.js ジョブキューライブラリ) の全ドキュメントを網羅したスキル。 ユーザーのタスクに応じて適切な README.md を読み、そこから個別ファイルへ辿ること。
ディレクトリ構造
.claude/skills/bullmq/
├── SKILL.md ← このファイル(エントリーポイント)
└── references/
├── start/README.md ← 入門ガイド(2ページ)
├── guide/README.md ← コアガイド(約50ページ)
│ ├── queues/README.md ← キュー関連(7ページ)
│ ├── workers/README.md ← ワーカー関連(8ページ)
│ ├── jobs/README.md ← ジョブ関連(13ページ)
│ ├── job-schedulers/README.md ← ジョブスケジューラ(4ページ)
│ ├── flows/README.md ← フロー関連(8ページ)
│ ├── events/README.md ← イベント(3ページ)
│ ├── metrics/README.md ← メトリクス(3ページ)
│ ├── telemetry/README.md ← テレメトリ(6ページ)
│ ├── redis-compatibility/README.md ← Redis互換性(3ページ)
│ ├── redis-hosting/README.md ← Redisホスティング(3ページ)
│ └── nestjs/README.md ← NestJS統合(3ページ)
├── patterns/README.md ← 実装パターン集(14ページ)
├── bullmq-pro/README.md ← BullMQ Pro(約18ページ)
│ ├── observables/README.md ← Observables(3ページ)
│ ├── groups/README.md ← Groups(10ページ)
│ └── nestjs/README.md ← NestJS Pro(3ページ)
└── languages/README.md ← 多言語サポート(3ページ)探索手順
1. ユーザーのタスクに最も関連するカテゴリを特定する 2. そのカテゴリの README.md を読む 3. README.md 内の一覧から必要な個別ファイルを選んで読む 4. 必要に応じて関連ページのリンクを辿る
カテゴリ → README.md マッピング
| タスク例 | カテゴリ | README パス |
|---|---|---|
| BullMQ の概要、インストール、基本的な使い方 | start | references/start/README.md |
| Queue, Worker, Job, Flow の基本操作 | guide | references/guide/README.md |
| キューの設定、バルク追加、グローバル並行数 | guide/queues | references/guide/queues/README.md |
| ワーカーの並行処理、シャットダウン、サンドボックス | guide/workers | references/guide/workers/README.md |
| ジョブの種類(遅延・繰り返し・優先度)、リトライ | guide/jobs | references/guide/jobs/README.md |
| ジョブスケジューラ、繰り返し戦略、cron | guide/job-schedulers | references/guide/job-schedulers/README.md |
| 親子ジョブ、FlowProducer、依存関係 | guide/flows | references/guide/flows/README.md |
| イベントリスニング、QueueEvents | guide/events | references/guide/events/README.md |
| メトリクス、Prometheus 連携 | guide/metrics | references/guide/metrics/README.md |
| OpenTelemetry、トレース、Jaeger | guide/telemetry | references/guide/telemetry/README.md |
| Redis 互換性、Dragonfly | guide/redis-compatibility | references/guide/redis-compatibility/README.md |
| AWS MemoryDB / ElastiCache | guide/redis-hosting | references/guide/redis-hosting/README.md |
| NestJS との統合 | guide/nestjs | references/guide/nestjs/README.md |
| 冪等性、スロットル、ステップ処理等の実装パターン | patterns | references/patterns/README.md |
| BullMQ Pro(グループ、Observable、バッチ) | bullmq-pro | references/bullmq-pro/README.md |
| Python / Elixir / PHP での利用 | languages | references/languages/README.md |
BullMQ Pro Batches
WorkerPro でジョブをバッチ処理する機能。複数のジョブを一度にまとめて処理することで、効率的なバルク操作が可能になる。
基本設定
batch オプションに size プロパティを渡してバッチ処理を有効にする:
import { WorkerPro } from '@taskforcesh/bullmq-pro';
const worker = new WorkerPro('myQueue', async (job) => {
const batch = job.getBatch();
// batch 内のジョブを処理
for (const batchedJob of batch) {
await processJob(batchedJob);
}
}, {
batch: {
size: 10, // 一度に最大10ジョブを処理
},
connection,
});高度なオプション
minSize と timeout
バッチ処理のタイミングを制御する追加パラメータ:
const worker = new WorkerPro('myQueue', processFn, {
batch: {
size: 50,
minSize: 10, // 最低10ジョブ集まるまで待機
timeout: 5000, // 最大5秒待機してから処理開始
},
connection,
});- minSize - 処理開始前に最低限必要なジョブ数
- timeout - minSize に達しない場合の最大待機時間(ミリ秒)。タイマー満了時に利用可能なジョブを処理する
ジョブの個別失敗処理
デフォルトでは例外がスローされるとバッチ内の全ジョブが失敗する。setAsFailed メソッドで個別にジョブを失敗させることが可能:
const worker = new WorkerPro('myQueue', async (job) => {
const batch = job.getBatch();
for (const batchedJob of batch) {
try {
await processJob(batchedJob);
} catch (err) {
batchedJob.setAsFailed(err);
}
}
}, {
batch: { size: 10 },
connection,
});イベント管理
バッチジョブは内部的にダミージョブでラップされる。Worker レベルのイベントリスナーはバッチコンテナのイベントを受信する:
worker.on('completed', async (job) => {
const batch = job.getBatch();
// batch 内の個別ジョブにアクセス
});個別ジョブのイベント監視には QueueEventsPro を使用する。
注意点
- バッチサイズは通常 10〜50 が推奨。大きなバッチはサイズに比例したオーバーヘッドが発生する
- minSize と timeout は Groups と互換性がない
- 以下の機能はバッチ処理で未サポート:
- 動的レート制限
- 手動ジョブ処理
- 動的ジョブ遅延
関連
- ./groups/README.md - Groups 機能
- ./observables/README.md - Observables 機能
BullMQ Pro Groups Concurrency
グループ単位で同時処理数を制限する機能。デフォルトではグループごとの並列実行に制限はないが、concurrency オプションで制御できる。
設定例
import { WorkerPro } from '@taskforcesh/bullmq-pro';
const worker = new WorkerPro('myQueue', processFn, {
group: {
concurrency: 3, // グループあたり最大3ジョブを並列処理
},
concurrency: 100, // Worker 全体の並行処理数
connection,
});この設定では、任意のグループに対して同時に3ジョブまでしか処理されない。Worker 全体の concurrency が100でも、各グループは最大3ジョブに制限される。
グローバルスコープ
グループの concurrency 設定はアプリケーション全体でグローバルに適用される。Worker インスタンスの数や個々の concurrency 設定に関わらず、1つのグループに対して同時に処理されるジョブ数は設定値を超えない。
例えば concurrency: 3 を設定した場合:
- Worker が1台でも10台でも、あるグループのジョブは最大3つまでしか同時実行されない
- 異なるグループのジョブは制限なく並列実行可能
注意点
- グループ concurrency はレート制限とは独立して動作する
- レート制限を設定していても、グループ concurrency を設定しないとグループ内で無制限に並列実行される
- Worker 全体の
concurrencyとグループのconcurrencyは別の概念 - ローカルグループ並行処理で個別のグループに異なる値を設定する場合は local-group-concurrency.md を参照
関連
- ./local-group-concurrency.md - ローカルグループ並行処理
- ./rate-limiting.md - グループ単位のレート制限
- ./groups.md - Groups の基本
BullMQ Pro Groups Getters
グループのジョブ数やジョブ一覧を取得するための API メソッド群。キューの状態監視やグループ管理に使用する。
グループ全体のジョブ数取得
getGroupsJobsCount() で全グループのジョブ数を取得する。優先度付きジョブと非優先度ジョブの両方が含まれる:
import { QueuePro } from '@taskforcesh/bullmq-pro';
const queue = new QueuePro('myQueue', { connection });
// イテレーションパラメータで取得範囲を指定
const count = await queue.getGroupsJobsCount(1000);
console.log('Total jobs across groups:', count);特定グループのアクティブジョブ数取得
getGroupActiveCount() で特定グループで現在処理中のジョブ数を取得する:
const groupId = 'myGroup';
const activeCount = await queue.getGroupActiveCount(groupId);
console.log(`Group ${groupId} active jobs:`, activeCount);グループ内のジョブ一覧取得
getGroupJobs() でページネーション形式でグループ内のジョブを取得する:
const groupId = 'myGroup';
const start = 0;
const end = 100;
const jobs = await queue.getGroupJobs(groupId, start, end);
console.log(`Jobs in group ${groupId}:`, jobs.length);注意点
getGroupsJobsCount()のイテレーションパラメータは取得するグループの走査範囲を指定するgetGroupJobs()はページネーション形式で、startとendで範囲を指定する- すべてのメソッドは
QueueProクラスで使用する
関連
- ./groups.md - Groups の基本的な使い方
- ./rate-limiting.md - グループ単位のレート制限
- ./concurrency.md - グループ単位の並行処理
BullMQ Pro Groups
1つのキュー内でジョブをグループに分配し、ラウンドロビン方式で公平に処理する機能。特定のユーザーやテナントがキューリソースを独占することを防ぐ。
概要
Groups は単一のキュー内に「仮想キュー」を作成する。例えば動画トランスコーディングサービスで多数のユーザーが1つのキューを共有する場合、1人のユーザーが大量のジョブを投入しても、グループ機能によりラウンドロビン方式で公平に処理される。
主な特徴
- 仮想キュー - ユーザーごとの仮想キューとして機能し、予測可能なジョブ処理を実現
- リソース効率 - 空のグループは Redis のリソースを消費しない
- 並列処理 - 複数 Worker や concurrency 設定時、グループ間でジョブが並列処理される
- 優先度 - グループに属さないジョブは、グループ内のジョブより優先される
- スケーラビリティ - Worker 数はキューの負荷に応じてスケール可能
ジョブの追加
group プロパティで id を指定してジョブをグループに追加する:
import { QueuePro } from '@taskforcesh/bullmq-pro';
const queue = new QueuePro('myQueue', { connection });
// グループ1にジョブを追加
const job1 = await queue.add(
'test',
{ foo: 'bar1' },
{
group: {
id: 1,
},
},
);
// グループ2にジョブを追加
const job2 = await queue.add(
'test',
{ foo: 'bar2' },
{
group: {
id: 2,
},
},
);グループジョブの処理
WorkerPro でジョブを処理する。job.opts.group でグループ情報にアクセス可能:
import { WorkerPro } from '@taskforcesh/bullmq-pro';
const worker = new WorkerPro('myQueue', async (job) => {
// 通常のジョブ処理
console.log('Processing job', job.id);
// グループ固有のロジック
await doSomethingSpecialForMyGroup(job.opts.group);
}, { connection });注意点
- グループ ID は数値または文字列を使用可能
- グループ数に実質的な上限はない
- グループに属さないジョブはグループジョブより優先して処理される
- 空のグループは自動的にクリーンアップされ、Redis リソースを消費しない
関連
- ./getters.md - グループ情報の取得
- ./rate-limiting.md - グループ単位のレート制限
- ./concurrency.md - グループ単位の並行処理
- ./pausing-groups.md - グループの一時停止
- ./prioritized.md - グループ内の優先度設定
BullMQ Pro Local Group Concurrency
個々のグループに固有の並行処理数を設定する機能。グループごとに異なる concurrency 値を持たせることができる。
concurrency の設定
setGroupConcurrency() メソッドで特定のグループに concurrency 値を割り当てる:
import { QueuePro } from '@taskforcesh/bullmq-pro';
const queue = new QueuePro('myQueue', { connection });
const groupId = 'my group';
await queue.setGroupConcurrency(groupId, 4);concurrency の取得
getGroupConcurrency() で特定グループの concurrency 設定を取得する:
const concurrency = await queue.getGroupConcurrency(groupId);
console.log(`Group concurrency: ${concurrency}`);Worker の設定
Worker インスタンスレベルでもグループ concurrency を設定する必要がある。これはデフォルト値として機能する:
import { WorkerPro } from '@taskforcesh/bullmq-pro';
const worker = new WorkerPro('myQueue', processFn, {
group: {
concurrency: 2, // デフォルトの並行処理数
},
concurrency: 100,
connection,
});注意点
- Worker インスタンスレベルでグループ concurrency を必ず設定すること。これが機能の前提条件であり、ローカル concurrency が未設定のグループのデフォルト値にもなる
- ローカル concurrency の値は Redis に保存されるため、不要になった場合は手動でクリーンアップが必要
- API リファレンス:
QueuePro.setGroupConcurrency(),QueuePro.getGroupConcurrency()
関連
- ./concurrency.md - グループ並行処理の基本
- ./local-group-rate-limit.md - ローカルグループレート制限
- ./groups.md - Groups の基本
BullMQ Pro Local Group Rate Limit
グループごとに個別のレート制限を設定する機能。ユーザーの課金プランやクォータに応じて、異なるレート制限を適用する場合に使用する。
概要
異なるグループに異なるレート制限を設定できる。デフォルトのレート制限は Worker で設定し、特定のグループには setGroupRateLimit() で個別の制限を上書きする。
実装例
import { QueuePro, WorkerPro } from '@taskforcesh/bullmq-pro';
const queue = new QueuePro('myQueue', { connection });
// 特定グループにローカルレート制限を設定
const groupId = 'my group';
const maxJobsPerDuration = 100;
const duration = 1000; // ミリ秒
await queue.setGroupRateLimit(groupId, maxJobsPerDuration, duration);
// Worker にデフォルトのレート制限を設定
const worker = new WorkerPro(
'myQueue',
async () => {
// ジョブ処理
},
{
group: {
limit: {
// デフォルトのレート制限設定
max: 1000,
duration: 1000,
},
},
connection,
},
);この例では「my group」グループに1秒あたり最大100ジョブのレート制限を設定している。他のグループにはデフォルトの1秒あたり1000ジョブが適用される。
注意点
- Worker インスタンスにデフォルトのレート制限(`group.limit`)を必ず設定すること。これがないとローカルレート制限が機能しない
setGroupRateLimit()は特定グループのデフォルト設定を上書きする- ローカルレート制限が設定されていないグループにはデフォルトのレート制限が適用される
- API リファレンス:
QueuePro.setGroupRateLimit()
関連
- ./rate-limiting.md - グループレート制限の基本
- ./local-group-concurrency.md - ローカルグループ並行処理
- ./groups.md - Groups の基本
BullMQ Pro Max Group Size
グループ内のジョブ数を制限する機能。maxSize オプションでグループの最大サイズを設定し、超過した場合は例外をスローする。
実装例
import { QueuePro, GroupMaxSizeExceededError } from '@taskforcesh/bullmq-pro';
const queue = new QueuePro('myQueue', { connection });
const groupId = 'my group';
try {
await queue.add('paint', { foo: 'bar' }, {
group: {
id: groupId,
maxSize: 7, // グループ内最大7ジョブ
},
});
} catch (err) {
if (err instanceof GroupMaxSizeExceededError) {
console.log(`Job discarded for group ${groupId}`);
} else {
throw err;
}
}動作
- グループが
maxSizeに達すると、新しいジョブの追加時にGroupMaxSizeExceededErrorがスローされる - エラーをキャッチして適切にハンドリングすることで、ジョブの破棄やリトライロジックを実装できる
注意点
maxSizeオプションはaddBulkメソッドでは使用できない(個別のジョブ追加のみ対応)- 超過分のジョブを破棄しても問題ない場合に使用する
maxSizeはジョブ追加時に毎回指定する必要がある
関連
- ./groups.md - Groups の基本
- ./concurrency.md - グループ単位の並行処理
- ./rate-limiting.md - グループ単位のレート制限
BullMQ Pro Pausing Groups
グループをグローバルに一時停止・再開する機能。一時停止されたグループのジョブは Worker によって取得されなくなる。
グループの一時停止
pauseGroup() メソッドでグループを一時停止する:
import { QueuePro } from '@taskforcesh/bullmq-pro';
const queue = new QueuePro('myQueue', { connection });
await queue.pauseGroup('groupId');- 一時停止されたグループのジョブは Worker に配信されない
- 現在処理中のジョブは完了するまで実行される(その後 Worker はアイドル状態になる)
- グループが既に一時停止中の場合は
falseを返す - 存在しないグループでも一時停止リストに追加される(エフェメラルグループに対応)
グループの再開
resumeGroup() メソッドでグループを再開する:
await queue.resumeGroup('groupId');- グループが存在しないか既に再開済みの場合は
falseを返す - 再開後、通常のジョブ処理が復帰する
使用例
import { QueuePro } from '@taskforcesh/bullmq-pro';
const queue = new QueuePro('myQueue', { connection });
// メンテナンスのためグループを一時停止
await queue.pauseGroup('user-123');
// メンテナンス完了後に再開
await queue.resumeGroup('user-123');注意点
- 一時停止はグローバルに適用され、すべての Worker に影響する
- 処理中のジョブは中断されず、完了まで実行される
- 存在しないグループを一時停止しても安全(後でグループが作成されたときに一時停止状態が適用される)
関連
- ./groups.md - Groups の基本
- ./rate-limiting.md - グループ単位のレート制限
- ./concurrency.md - グループ単位の並行処理
BullMQ Pro Prioritized Intra-Groups
グループ内でジョブに優先度を割り当てる機能。group オプションと priority オプションを同時に指定することで有効になる。
優先度付きジョブの追加
import { QueuePro } from '@taskforcesh/bullmq-pro';
const queue = new QueuePro('myQueue', { connection });
await queue.add(
'paint',
{ foo: 'bar' },
{
group: {
id: 'groupId',
priority: 10,
},
},
);優先度の範囲
- 優先度は 0 から 2097151 の範囲で指定する
- 数値が大きいほど優先度が低い(Unix プロセスと同じ方式)
- 優先度を明示的に指定しないジョブはデフォルトで最高優先度(0)が割り当てられる
優先度ごとのジョブ数取得
getCountsPerPriorityForGroup() でグループ内の優先度別ジョブ数を取得する:
const counts = await queue.getCountsPerPriorityForGroup('groupId', [1, 0]);
/*
{
'1': 11,
'0': 10
}
*/この例ではグループ「groupId」内の優先度1のジョブが11個、優先度0のジョブが10個あることを示している。
注意点
- 優先度は
groupオプション内に指定する(トップレベルのpriorityとは異なる) - 優先度0が最高優先度
- 数値が大きいほど後回しにされる
- グループに属さないジョブの優先度とは独立して動作する
関連
- ./groups.md - Groups の基本
- ./concurrency.md - グループ単位の並行処理
- ./rate-limiting.md - グループ単位のレート制限
BullMQ Pro Groups Rate Limiting
グループ単位で独立したレート制限を適用する機能。各グループが時間単位あたりの最大ジョブ数を超えると、そのグループのみが制限される。
静的レート制限設定
Worker インスタンスでグループのレート制限を設定する:
import { WorkerPro } from '@taskforcesh/bullmq-pro';
const worker = new WorkerPro('myQueue', processFn, {
group: {
limit: {
max: 100, // グループあたり最大100ジョブ
duration: 1000, // 1秒あたり
},
},
connection,
});- max - 時間枠内でグループあたりに許可される最大ジョブ数
- duration - 時間枠(ミリ秒)
手動レート制限
外部 API のレスポンスなど動的な条件に基づいてレート制限を適用する:
import { WorkerPro, Worker } from '@taskforcesh/bullmq-pro';
const worker = new WorkerPro(
'myQueue',
async (job) => {
const groupId = job.opts.group.id;
const [isRateLimited, duration] = await doExternalCall(groupId);
if (isRateLimited) {
await worker.rateLimitGroup(job, duration);
throw Worker.RateLimitError();
}
},
{ connection },
);レート制限状態の確認
getGroupRateLimitTtl() でグループが現在レート制限中かどうかを確認する:
import { QueuePro } from '@taskforcesh/bullmq-pro';
const queue = new QueuePro('myQueue', { connection });
const groupId = '0';
const maxJobs = 100;
const ttl = await queue.getGroupRateLimitTtl(groupId, maxJobs);
if (ttl > 0) {
console.log('Group is rate limited, remaining TTL:', ttl);
}TTL が正の値の場合、そのグループは現在レート制限中。
注意点
rateLimitGroup()呼び出し後、ジョブはアクティブ状態ではなくなるため、必ずWorker.RateLimitError()をスローすること- レート制限はグループごとに独立して適用される
- 他のグループのジョブ処理には影響しない
- ローカルグループレート制限との併用については local-group-rate-limit.md を参照
関連
- ./local-group-rate-limit.md - ローカルグループレート制限
- ./concurrency.md - グループ単位の並行処理
- ./groups.md - Groups の基本
BullMQ Pro — Groups
ジョブをグループ化し、グループごとに独立した並行処理・レート制限を適用する。
| Name | Description | Path |
|---|---|---|
| Groups | グループの概要と基本的な使い方 | ./groups.md |
| Getters | グループ情報の取得 | ./getters.md |
| Rate Limiting | グループ単位のレート制限 | ./rate-limiting.md |
| Local Group Rate Limit | ローカルグループレート制限 | ./local-group-rate-limit.md |
| Concurrency | グループ単位の並行処理 | ./concurrency.md |
| Local Group Concurrency | ローカルグループ並行処理 | ./local-group-concurrency.md |
| Max Group Size | グループの最大サイズ | ./max-group-size.md |
| Pausing Groups | グループの一時停止 | ./pausing-groups.md |
| Prioritized | グループ内の優先度設定 | ./prioritized.md |
| Sandboxes for Groups | グループ用サンドボックスプロセッサ | ./sandboxes-for-groups.md |
BullMQ Pro Sandboxes for Groups
サンドボックスプロセッサでグループ処理を使用する機能。サンドボックス内のジョブオブジェクトから gid プロパティでグループ ID にアクセスできる。
実装例
import { SandboxedJobPro } from '@taskforcesh/bullmq-pro';
module.exports = function (job: SandboxedJobPro) {
// グループ ID にアクセス
const groupId = job.gid;
console.log('Processing job for group:', groupId);
// job.opts.group からもグループ情報にアクセス可能
console.log('Group options:', job.opts.group);
console.log('Group ID:', job.opts.group.id);
// gid と opts.group.id は同じ値
// job.gid === job.opts.group.id
};プロパティ
サンドボックスプロセッサ内で利用可能なグループ関連プロパティ:
job.gid- グループ ID(文字列)job.opts.group- グループオプションオブジェクトjob.opts.group.id- グループ ID(job.gidと同値)
注意点
- サンドボックスプロセッサで対応している Pro 機能は Groups のみ。Observables やその他の Pro 機能はサンドボックス環境では利用不可
SandboxedJobPro型を使用すること(通常のSandboxedJobではなく)gidは常に文字列型として提供される
関連
- ./groups.md - Groups の基本
- ./concurrency.md - グループ単位の並行処理
BullMQ Pro Install
BullMQ Pro のインストールには taskforce.sh から取得する NPM トークンが必要。.npmrc ファイルでレジストリ設定を行い、パッケージをインストールする。
.npmrc の設定
リポジトリのルートに .npmrc ファイルを作成または更新する:
@taskforcesh:registry=https://npm.taskforce.sh/
//npm.taskforce.sh/:_authToken=${NPM_TASKFORCESH_TOKEN}
always-auth=trueインストール
パッケージマネージャーでインストールする:
yarn add @taskforcesh/bullmq-proまたは npm を使用:
npm install @taskforcesh/bullmq-pro基本的な使い方
Pro 版のクラスをインポートして使用する:
import { QueuePro, WorkerPro } from '@taskforcesh/bullmq-pro';
const queue = new QueuePro('myQueue', {
connection: {
host: 'localhost',
port: 6379,
},
});
const worker = new WorkerPro('myQueue', async (job) => {
// ジョブを処理
}, {
connection: {
host: 'localhost',
port: 6379,
},
});Docker での設定
Docker ビルドで .npmrc ファイルを含める:
WORKDIR /app
ADD .npmrc /app/.npmrc
RUN npm install注意点
NPM_TASKFORCESH_TOKEN環境変数にトークンを設定する必要がある.npmrcファイルにトークンを直接ハードコードしないこと- CI/CD 環境ではシークレット変数としてトークンを管理する
- Pro 版のクラス名は
QueuePro、WorkerPro、FlowProducerProなど、末尾にProが付く
関連
- ./introduction.md - BullMQ Pro の概要
- ./support.md - サポートとライセンス
BullMQ Pro Introduction
BullMQ Pro は BullMQ の商用版であり、オープンソース版を超える高度な機能と、ライブラリ作者による専用サポートを提供する。
概要
BullMQ Pro は標準の BullMQ を直接置き換える形で動作し、大きな統合変更なしにプレミアム機能を利用できる。
主な機能
BullMQ Pro は以下の高度な機能を提供する:
- Groups - ジョブをグループ化し、グループごとにラウンドロビンで処理
- Observables - Observable パターンによるジョブ処理、キャンセル、TTL
- Batches - 複数ジョブの一括処理
- Telemetry - OpenTelemetry によるモニタリング
ライセンスモデル
BullMQ Pro は組織単位のライセンスで、すべてのプロジェクトで無制限に使用可能。購入前に無料トライアル期間が提供される。
継続的な開発
新しい機能が定期的に追加されており、プロジェクトのロードマップで開発の進捗を確認できる。
注意点
- BullMQ Pro は有料ライセンスが必要
- npm レジストリトークンが必要(インストール手順は install.md を参照)
- オープンソース版の BullMQ とは別パッケージ(
@taskforcesh/bullmq-pro)
関連
- ./install.md - インストール方法
- ./observables/README.md - Observables 機能
- ./groups/README.md - Groups 機能
- ./batches.md - Batches 機能
BullMQ Pro NestJS Producers
NestJS で BullMQ Pro のキュープロデューサーとフロープロデューサーを使用する方法。@InjectQueue と @InjectFlowProducer デコレータで DI を通じてキューにアクセスする。
キュープロデューサー
@InjectQueue() デコレータでキューを注入し、ジョブを追加する:
import { Injectable } from '@nestjs/common';
import { QueuePro } from '@taskforcesh/bullmq-pro';
import { InjectQueue } from '@taskforcesh/nestjs-bullmq-pro';
@Injectable()
export class AudioService {
constructor(@InjectQueue('audio') private audioQueue: QueuePro) {}
async addJob() {
const job = await this.audioQueue.add('transcode', {
foo: 'bar',
});
return job;
}
}@InjectQueue()デコレータのパラメータはregisterQueue()で指定したキュー名と一致させる- 注入される型は
QueuePro(通常のQueueではなく)
フロープロデューサー
親子関係を持つ複雑なジョブワークフローを作成する:
import { Injectable } from '@nestjs/common';
import { FlowProducerPro } from '@taskforcesh/bullmq-pro';
import { InjectFlowProducer } from '@taskforcesh/nestjs-bullmq-pro';
@Injectable()
export class FlowService {
constructor(
@InjectFlowProducer('flow') private fooFlowProducer: FlowProducerPro,
) {}
async addFlow() {
const job = await this.fooFlowProducer.add({
name: 'root-job',
queueName: 'topQueueName',
data: {},
children: [
{
name: 'child-job',
data: { idx: 0, foo: 'bar' },
queueName: 'childrenQueueName',
},
],
});
return job;
}
}@InjectFlowProducer()デコレータで登録済みのフロー名を指定する- フロー構造はルートジョブと子ジョブの階層で定義する
- 複数キューにまたがるジョブ定義が可能
注意点
@taskforcesh/nestjs-bullmq-proパッケージを別途インストールする必要がある- Pro 版のデコレータとクラスを使用すること(
@InjectQueueは@taskforcesh/nestjs-bullmq-proから、型はQueuePro/FlowProducerPro) - キュー名はモジュールの
registerQueue()で登録した名前と一致させる
関連
- ./queue-events-listeners.md - イベントリスナー
- ../groups/README.md - Groups 機能
- ../batches.md - バッチ処理
BullMQ Pro NestJS Queue Events Listeners
NestJS で BullMQ Pro のキューイベントをデコレータベースでリッスンする方法。@QueueEventsListener と @OnQueueEvent デコレータでジョブのライフサイクルイベントに応答する。
イベントリスナークラスの作成
@QueueEventsListener デコレータで監視対象のキューを指定し、QueueEventsHost を継承する:
import {
QueueEventsListener,
QueueEventsHost,
OnQueueEvent,
} from '@taskforcesh/nestjs-bullmq-pro';
@QueueEventsListener('queueName')
export class TestQueueEvents extends QueueEventsHost {
@OnQueueEvent('completed')
onCompleted({
jobId,
}: {
jobId: string;
returnvalue: string;
prev?: string;
}) {
console.log(`Job ${jobId} completed`);
}
}モジュールへの登録
イベントリスナークラスをモジュールの providers に登録する:
import { Module } from '@nestjs/common';
import { BullModule } from '@taskforcesh/nestjs-bullmq-pro';
@Module({
imports: [
BullModule.registerQueue({
name: 'queueName',
connection: {
host: '0.0.0.0',
port: 6380,
},
}),
],
providers: [TestQueueEvents],
})
export class AppModule {}主要なイベント
@OnQueueEvent デコレータで監視可能な主なイベント:
'completed'- ジョブ完了時'failed'- ジョブ失敗時'progress'- ジョブ進捗更新時'active'- ジョブがアクティブになった時'waiting'- ジョブが待機状態になった時
注意点
@QueueEventsListenerのパラメータはregisterQueue()で登録したキュー名と一致させる- イベントハンドラは
jobId、returnvalue、prevなどのメタデータを含むオブジェクトを受け取る - NestJS の DI システムとシームレスに統合される
@taskforcesh/nestjs-bullmq-proパッケージのデコレータを使用すること
関連
- ./producers.md - キュープロデューサー
- ../groups/README.md - Groups 機能
BullMQ Pro — NestJS
BullMQ Pro の NestJS 統合モジュール。Pro 版のキュー、ワーカー、イベントリスナーを NestJS のデコレータと DI で利用する。
| Name | Description | Path |
|---|---|---|
| Producers | NestJS Pro でのキュープロデューサー | ./producers.md |
| Queue Events Listeners | NestJS Pro でのイベントリスナー | ./queue-events-listeners.md |
BullMQ Pro Observable Cancelation
Observable ジョブに TTL(Time to Live)を設定し、処理時間が長すぎるジョブを自動的にキャンセルする機能。
概要
BullMQ Pro では TTL 値を設定することで、ジョブの最大処理時間を定義できる。処理時間が TTL を超えた場合、Observable が自動的にキャンセル(unsubscribe)される。
グローバル TTL 設定
すべてのジョブに統一的な TTL を設定する:
import { WorkerPro } from '@taskforcesh/bullmq-pro';
const worker = new WorkerPro(queueName, processor, {
ttl: 100, // 全ジョブに100msのTTL
connection,
});ジョブ名ごとの TTL 設定
ジョブ名に応じて異なる TTL を設定する:
import { WorkerPro } from '@taskforcesh/bullmq-pro';
const worker = new WorkerPro(queueName, processor, {
ttl: {
test1: 100, // test1 ジョブは100ms
test2: 200, // test2 ジョブは200ms
},
connection,
});注意点
- TTL はミリ秒単位で指定する
- TTL を超えた場合、Observable の
unsubscribe関数が呼ばれるため、リソースのクリーンアップが可能 - グローバル TTL とジョブ名ごとの TTL は排他的(どちらか一方を指定)
- TTL は Observable を返すプロセッサでのみ有効
関連
- ./observables.md - Observable の基本的な使い方
- ../groups/README.md - Groups 機能
BullMQ Pro Observables
Worker から Promise の代わりに Observable を返すことで、複数値の送出、ジョブのキャンセル、状態復元などの高度なユースケースに対応する。
概要
Observables は Promise に対して以下の2つの主要な利点を持つ:
- 複数の値を送出可能 -
subscriber.next()で複数回値を返せる - キャンセル可能 - 実行中のジョブをクリーンにキャンセルできる
主なユースケース
- ジョブのキャンセル - Observable を通じてクリーンな終了処理
- TTL(Time to Live) - 処理時間が長すぎるジョブの自動キャンセル
- 状態の永続化とリトライ - Observable が返す最後の値が永続化されるため、リトライ時に前回の続きから再開可能
基本的な使い方
import { WorkerPro } from '@taskforcesh/bullmq-pro';
import { Observable } from 'rxjs';
const processor = async () => {
return new Observable<number>((subscriber) => {
subscriber.next(1);
subscriber.next(2);
subscriber.next(3);
const intervalId = setTimeout(() => {
subscriber.next(4);
subscriber.complete();
}, 500);
// キャンセル時のクリーンアップ関数
return function unsubscribe() {
clearInterval(intervalId);
};
});
};
const worker = new WorkerPro(queueName, processor, { connection });この例では Observable が4つの値を送出する。最初の3つは即座に、4つ目は500ms後に送出される。subscriber が返す unsubscribe 関数は Observable がキャンセルされた際に呼ばれるクリーンアップ処理。
状態復元パターン(ステートマシン)
クラッシュからの復帰シナリオで、最後に永続化された値から処理を再開する:
import { WorkerPro } from '@taskforcesh/bullmq-pro';
import { Observable } from 'rxjs';
const processor = async (job) => {
return new Observable<number>((subscriber) => {
switch (job.returnvalue) {
default:
subscriber.next(1);
// fall through
case 1:
subscriber.next(2);
// fall through
case 2:
subscriber.next(3);
// fall through
case 3:
subscriber.complete();
}
});
};
const worker = new WorkerPro(queueName, processor, { connection });job.returnvalue には前回の実行で最後に送出された値が格納されている。switch 文の fall-through を利用して、前回の続きから処理を再開する。
注意点
rxjsパッケージを別途インストールする必要がある- Observable の最後の値が
job.returnvalueとして永続化される subscriber.complete()を呼ばないと Observable は完了しないunsubscribe関数でリソースのクリーンアップを必ず行うこと
関連
- ./cancelation.md - Observable ジョブのキャンセル
- ../batches.md - バッチ処理
BullMQ Pro — Observables
Observable パターンによるジョブの進捗ストリーミング、キャンセル、状態復元を提供する。
| Name | Description | Path |
|---|---|---|
| Observables | Observable パターンによるジョブの進捗ストリーミング | ./observables.md |
| Cancelation | Observable ジョブのキャンセル | ./cancelation.md |
BullMQ Pro
BullMQ の有料拡張版。Groups、Observables、Batches などの高度な機能を提供する。
概要
| Name | Description | Path |
|---|---|---|
| Introduction | BullMQ Pro の概要と機能一覧 | ./introduction.md |
| Install | BullMQ Pro のインストール方法 | ./install.md |
機能
| Name | Description | Path |
|---|---|---|
| Observables | Observable パターンによるジョブ処理 | ./observables/README.md |
| Groups | ジョブのグループ化と独立した制御 | ./groups/README.md |
| Batches | ジョブのバッチ処理 | ./batches.md |
| Telemetry | BullMQ Pro のテレメトリ | ./telemetry.md |
統合
| Name | Description | Path |
|---|---|---|
| NestJS | NestJS との統合(Pro 版) | ./nestjs/README.md |
| Support | サポートとライセンス | ./support.md |
BullMQ Pro Support
BullMQ Pro には商用サポートが含まれており、メールベースでライブラリ作者から直接サポートを受けられる。
サポート方法
- メールサポート: support@taskforce.sh にメールで問い合わせや問題を報告できる
レスポンス時間
- 初回応答は 1 営業日以内を目標としている
- ベストエフォートでより迅速な対応が行われる場合もある
- 解決時間は問題の複雑さに依存し、単純な問い合わせは迅速に処理されるが、複雑な技術的問題は数日以上かかる場合がある
- 完全な修正ではなく回避策が提供される場合もある
ライセンス情報
- サブスクリプションの詳細はアカウントポータルの BullMQ Pro タブで確認できる
- 組織単位のライセンスで、すべてのプロジェクトで無制限に使用可能
注意点
- サポートは有料ライセンス保有者のみが利用可能
- レスポンス時間は目安であり、保証ではない
- 最新のサブスクリプション情報はアカウントポータルで確認すること
関連
- ./introduction.md - BullMQ Pro の概要
- ./install.md - インストール方法
BullMQ Pro Telemetry
BullMQ Pro は OpenTelemetry を活用したテレメトリ機能をサポートする。オープンソース版と同じ統合方式で、Pro 固有の機能(Groups、Batches)の監視も可能。
Queue の設定
QueuePro にテレメトリプロバイダーを設定する:
import { QueuePro } from '@taskforcesh/bullmq-pro';
import { BullMQOtel } from 'bullmq-otel';
const queue = new QueuePro('myProQueue', {
connection,
telemetry: new BullMQOtel('guide'),
});ジョブの追加(グループ付き)
テレメトリが有効な Queue にグループ付きジョブを追加する:
await queue.add(
'myJob',
{ data: 'myData' },
{
attempts: 2,
backoff: 1000,
group: {
id: 'myGroupId',
},
},
);Worker の設定
Worker にもテレメトリインスタンスを渡す:
import { WorkerPro } from '@taskforcesh/bullmq-pro';
import { BullMQOtel } from 'bullmq-otel';
const worker = new WorkerPro(
'myProQueue',
async (job) => {
console.log('processing job', job.id);
},
{
name: 'myWorker',
connection,
telemetry: new BullMQOtel('guide'),
concurrency: 10,
batch: { size: 10 },
},
);注意点
bullmq-otelパッケージを別途インストールする必要がある- テレメトリは Queue と Worker の両方に設定することを推奨
- Pro 固有機能(Groups、Batches)のスパンも自動的に収集される
- OpenTelemetry の詳細な統合方法については公式ブログの包括的なチュートリアルを参照
関連
- ./install.md - インストール方法
- ./batches.md - バッチ処理
- ./groups/README.md - Groups 機能
Architecture
BullMQ は Redis 上にジョブキュー機能を実装しており、明確に定義されたライフサイクルステート管理システムを持っています。ジョブは追加時に初期状態に入り、処理を経て完了または失敗状態に遷移します。
ジョブのライフサイクル
ジョブが Queue.add() で追加されると、以下のいずれかの初期状態に入ります。
| State | Description |
|---|---|
wait | 処理前の標準待機リスト |
prioritized | 優先度順に並べられたジョブ(0 が最高優先度、Unix nice 規約に従う) |
delayed | 将来の処理のためにスケジュールされたジョブ。時間が来ると wait または prioritized に移動 |
処理フロー
add() → wait/prioritized/delayed
↓
active(処理中)
↓
completed / failedステート遷移
// ジョブの追加(wait 状態へ)
const queue = new Queue('myQueue');
await queue.add('myJob', { data: 'value' });
// 遅延ジョブの追加(delayed 状態へ)
await queue.add('delayedJob', { data: 'value' }, {
delay: 5000, // 5秒後に wait に移動
});
// 優先度付きジョブの追加(prioritized 状態へ)
await queue.add('priorityJob', { data: 'value' }, {
priority: 1, // 高優先度
});| State | Description |
|---|---|
active | ワーカーが処理中の状態。ロックが取得される |
completed | 処理が正常に完了した状態 |
failed | プロセッサが例外を throw した、またはストールした状態 |
waiting-children | 子ジョブの完了を待っている親ジョブの状態(FlowProducer 使用時) |
FlowProducer パターン
FlowProducer を使用すると、親子関係を持つジョブフローを作成できます。親ジョブは waiting-children 状態に入り、すべての子ジョブが完了すると wait 状態に移動して処理されます。
import { FlowProducer } from 'bullmq';
const flowProducer = new FlowProducer();
await flowProducer.add({
name: 'parentJob',
queueName: 'parentQueue',
data: {},
children: [
{
name: 'childJob1',
queueName: 'childQueue',
data: { step: 1 },
},
{
name: 'childJob2',
queueName: 'childQueue',
data: { step: 2 },
},
],
});完全なライフサイクル図
┌───────────┐
│ delayed │
└─────┬─────┘
│ (時間経過)
▼
add() ──→ ┌───────────────────────────┐
│ wait / prioritized │
└─────────────┬─────────────┘
│
▼
┌───────────┐
│ active │
└─────┬─────┘
│
┌─────────┴─────────┐
▼ ▼
┌───────────┐ ┌───────────┐
│ completed │ │ failed │
└───────────┘ └───────────┘
FlowProducer:
add() ──→ waiting-children ──→ wait ──→ active ──→ completed/failed注意点
- Redis はすべてのジョブデータとステート管理のバックエンドとして使用されます
- アクティブなジョブにはロックが設定され、他のワーカーが同じジョブを処理することを防ぎます
- ロックの有効期限はデフォルトで 30 秒で、ワーカーが定期的に更新します
maxmemory-policyは必ずnoevictionに設定する必要があります
関連
- ./going-to-production.md
- ./parallelism-and-concurrency.md
- ./troubleshooting.md
BullMQ — Connections
BullMQ は Redis への接続に ioredis モジュールを使用する。Queue や Worker の各インスタンスに接続設定を渡す方法、共有接続の注意点、および必須の Redis 設定について説明する。
基本的な接続
接続オプションを指定しない場合、デフォルトで localhost:6379 に接続する。
import { Queue, Worker } from 'bullmq';
// 各インスタンスに個別の接続設定を渡す
const myQueue = new Queue('myqueue', {
connection: {
host: 'myredis.taskforce.run',
port: 32856,
},
});
const myWorker = new Worker('myqueue', async job => {}, {
connection: {
host: 'myredis.taskforce.run',
port: 32856,
},
});接続の共有(Reusing Connections)
Producer 間の共有
Queue(Producer)同士は ioredis インスタンスを共有できる。
import { Queue } from 'bullmq';
import IORedis from 'ioredis';
const connection = new IORedis();
// 2つの Producer で同じ接続を再利用
const myFirstQueue = new Queue('myFirstQueue', { connection });
const mySecondQueue = new Queue('mySecondQueue', { connection });Consumer 間の共有
Worker(Consumer)間で接続を共有する場合、maxRetriesPerRequest: null の設定が必須。
import { Worker } from 'bullmq';
import IORedis from 'ioredis';
const connection = new IORedis({ maxRetriesPerRequest: null });
// 2つの Consumer で同じ接続を再利用
const myFirstWorker = new Worker('myFirstWorker', async job => {}, {
connection,
});
const mySecondWorker = new Worker('mySecondWorker', async job => {}, {
connection,
});接続オプション
| Name | Type | Description |
|---|---|---|
host | string | Redis ホスト名(デフォルト: localhost) |
port | number | Redis ポート番号(デフォルト: 6379) |
password | string | Redis 認証パスワード |
db | number | Redis データベース番号 |
maxRetriesPerRequest | `number \ | null` |
prefix | string | BullMQ キーのプレフィックス |
Producer と Consumer の違い
- Producer(HTTP エンドポイントからジョブを追加する場合): デフォルトのリトライ設定、または
maxRetriesPerRequest: 1のように低い値を使用し、Redis 不通時に速やかに失敗させる - Consumer(バックグラウンドワーカー): 非同期で動作するため、より多くのリトライ回数を許容できる。
maxRetriesPerRequest: nullが推奨
注意点
- keyPrefix を使用しないこと: ioredis の
keyPrefixオプションは使用禁止。BullMQ は独自のキー・プレフィックス機構(prefixオプション)を提供している - maxmemory-policy: Redis の設定で
maxmemory-policy=noevictionが必須。これを設定しないと、Redis がキーを自動削除し BullMQ の動作が壊れる可能性がある - QueueScheduler / QueueEvents: これらはブロッキング接続を必要とするため、接続の共有ができない
関連
- Introduction
- Queues
- Workers
Create Custom Events
QueueEventsProducer クラスを使用して、BullMQ 上で汎用的な分散リアルタイムイベントエミッターを作成できます。コンシューマーは QueueEvents クラスでイベントをサブスクライブします。
基本的な使い方
const queueName = 'customQueue';
const queueEventsProducer = new QueueEventsProducer(queueName, {
connection,
});
const queueEvents = new QueueEvents(queueName, {
connection,
});
interface CustomListener extends QueueEventsListener {
example: (args: { custom: string }, id: string) => void;
}
queueEvents.on<CustomListener>('example', async ({ custom }) => {
// custom logic
});
interface CustomEventPayload {
eventName: string;
custom: string;
}
await queueEventsProducer.publishEvent<CustomEventPayload>({
eventName: 'example',
custom: 'value',
});QueueEventsProducer API
| Name | Type | Description |
|---|---|---|
| queueName | string | イベントを発行するキュー名 |
| connection | ConnectionOptions | Redis 接続設定 |
| publishEvent | method | カスタムイベントを発行するメソッド |
イベントペイロード
| Name | Type | Description |
|---|---|---|
| eventName | string | イベント名(必須) |
| その他のプロパティ | any | カスタムデータ |
注意点
eventName属性のみが必須です- 一部のイベント名は予約されています(QueueListener API Reference を参照)
- カスタムリスナーのインターフェースを定義して型安全にイベントを扱えます
関連
- ./events.md — Worker イベントと QueueEvents の基本
Events
BullMQ のすべてのクラスは EventEmitter を継承しており、ジョブのライフサイクルに関するイベントを発行します。Worker のローカルイベントと、QueueEvents クラスによるグローバル監視の2種類の方法でイベントをリスンできます。
ローカル Worker イベント
Worker はローカルで処理したジョブに関するイベントを発行します:
import { Queue } from 'bullmq';
const myQueue = new Queue('Paint');
myQueue.on('waiting', (job: Job) => {
// Job is waiting to be processed.
});import { Worker } from 'bullmq';
const myWorker = new Worker('Paint');
myWorker.on('drained', () => {
// Queue is drained, no more jobs left
});
myWorker.on('completed', (job: Job) => {
// job has completed
});
myWorker.on('failed', (job: Job) => {
// job has failed
});グローバルイベント監視(QueueEvents)
すべての Worker からのイベントを一箇所でリスンするには QueueEvents クラスを使用します:
import { QueueEvents } from 'bullmq';
const queueEvents = new QueueEvents('Paint');
queueEvents.on('completed', ({ jobId }) => {
// Called every time a job is completed in any worker.
});
queueEvents.on(
'progress',
({ jobId, data }: { jobId: string; data: number | object }) => {
// jobId received a progress event
},
);主な Worker イベント
| Event | Description |
|---|---|
| waiting | ジョブが処理待ちキューに入った |
| completed | ジョブが正常に完了した |
| failed | ジョブがエラーで失敗した |
| progress | ジョブの進捗が更新された |
| drained | キューにジョブがなくなった |
| error | Worker でエラーが発生した |
Redis Streams ベースの実装
QueueEvents は Redis Streams を使用して実装されており、標準的な pub-sub と異なり、ネットワーク切断時にもイベントが失われない保証があります。
イベントストリームは自動的にトリミングされ、デフォルトで約10,000件のイベントが保持されます。この上限は streams.events.maxLen オプションで変更できます。
手動イベントトリミング
import { Queue } from 'bullmq';
const queue = new Queue('paint');
await queue.trimEvents(10); // 最新10件のイベントのみ残す注意点
- ローカルイベントは各 Worker が処理したジョブのみ対象です
- すべての Worker のイベントを一元監視する場合は
QueueEventsを使用してください - イベントストリームのサイズはデフォルトで約10,000件に自動トリミングされます
trimEvents()で手動トリミングも可能です
関連
- ./create-custom-events.md — カスタムイベントの作成
BullMQ — Events
| Name | Description | Path |
|---|---|---|
| Events | Worker イベントと QueueEvents によるイベント監視 | ./events.md |
| Create Custom Events | カスタムイベントの作成と発行 | ./create-custom-events.md |
Adding Flows in Bulk
FlowProducer.addBulk() を使用して、複数のフローをアトミックに一括追加できます。すべてのフローが作成されるか、いずれも作成されないかのどちらかとなり、Redis へのラウンドトリップ回数も削減されます。
基本的な使い方
import { FlowProducer } from 'bullmq';
const flow = new FlowProducer({ connection });
const trees = await flow.addBulk([
{
name: 'root-job-1',
queueName: 'rootQueueName-1',
data: {},
children: [
{
name,
data: { idx: 0, foo: 'bar' },
queueName: 'childrenQueueName-1',
},
],
},
{
name: 'root-job-2',
queueName: 'rootQueueName-2',
data: {},
children: [
{
name,
data: { idx: 1, foo: 'baz' },
queueName: 'childrenQueueName-2',
},
],
},
]);注意点
- この呼び出しは成功か失敗のどちらかのみで、すべてのジョブが追加されるか、まったく追加されないかのいずれかです
- 大量のフローを追加する場合、個別の
add()呼び出しよりも高速です
関連
- ./flows.md — FlowProducer の基本
- ./get-flow-tree.md — フローツリーの取得
Continue Parent
continueParentOnFailure オプションにより、子ジョブが失敗した場合に親ジョブを即座に処理開始できます。removeUnprocessedChildren メソッドで未処理の子ジョブを動的にクリーンアップし、getFailedChildrenValues() で失敗原因を判別できます。(v5.58.0 以降)
continueParentOnFailure
const { FlowProducer } = require('bullmq');
const flow = new FlowProducer({ connection });
const originalTree = await flow.add({
name: 'root-job',
queueName: 'topQueueName',
data: {},
children: [
{
name: 'child-job-1',
data: { idx: 0, foo: 'bar' },
queueName: 'childrenQueueName',
opts: { continueParentOnFailure: true },
},
{
name: 'child-job-2',
data: { idx: 1, foo: 'baz' },
queueName: 'childrenQueueName',
},
{
name: 'child-job-3',
data: { idx: 2, foo: 'qux' },
queueName: 'childrenQueueName',
},
],
});子ジョブに continueParentOnFailure: true を設定すると、その子が失敗した時点で親ジョブが即座に active 状態に移行します(他の子ジョブがまだ処理中でも)。
removeUnprocessedChildren
未処理(waiting または delayed 状態)の子ジョブをすべて削除します:
await job.removeUnprocessedChildren();| 対象 | 説明 |
|---|---|
| 削除される | waiting, delayed 状態の子ジョブ |
| 削除されない | active, completed, failed 状態の子ジョブ |
getFailedChildrenValues
失敗した子ジョブの ID とエラーメッセージのマッピングを取得します:
const failedChildren = await job.getFailedChildrenValues();
// { "job-id-1": "Upload failed" }失敗した子がない場合は空オブジェクトを返します。
親 Worker での分岐処理
const processor = async (job) => {
const failedChildren = await job.getFailedChildrenValues();
const hasFailedChildren = Object.keys(failedChildren).length > 0;
if (hasFailedChildren) {
// パス1: 子ジョブの失敗により continueParentOnFailure がトリガー
console.log(`Parent job ${job.name} triggered by child failure(s):`, failedChildren);
// 未処理の子ジョブを削除
await job.removeUnprocessedChildren();
console.log('Unprocessed child jobs have been removed.');
} else {
// パス2: すべての子ジョブが正常完了
console.log(`Parent job ${job.name} processing after all children completed successfully.`);
}
};注意点
continueParentOnFailureは子ジョブの失敗時に親を即座に active 状態にします(デフォルトでは全子ジョブの完了を待つ)removeUnprocessedChildrenは active, completed, failed 状態のジョブには影響しません- ファイルアップロードなど、1つの失敗で残りを中止すべきワークフローに適しています
関連
- ./flows.md — FlowProducer の基本
- ./fail-parent.md — 子失敗時の親の失敗
- ./remove-dependency.md — 依存関係の削除
Fail Parent
failParentOnFailure オプションを子ジョブに設定すると、その子ジョブが失敗した場合に親ジョブも失敗としてマークされます。この効果はジョブ階層を再帰的に伝播し、祖父母以上のジョブも失敗させることができます。
基本的な使い方
import { FlowProducer } from 'bullmq';
const flow = new FlowProducer({ connection });
const originalTree = await flow.add({
name: 'root-job',
queueName: 'topQueueName',
data: {},
children: [
{
name: 'child-job',
data: { idx: 0, foo: 'bar' },
queueName: 'childrenQueueName',
opts: { failParentOnFailure: true },
children: [
{
name,
data: { idx: 1, foo: 'bah' },
queueName: 'grandChildrenQueueName',
opts: { failParentOnFailure: true },
},
{
name,
data: { idx: 2, foo: 'baz' },
queueName: 'grandChildrenQueueName',
// failParentOnFailure なし: この子の失敗は親に影響しない
},
],
},
{
name,
data: { idx: 3, foo: 'foo' },
queueName: 'childrenQueueName',
// failParentOnFailure なし: この子の失敗は親に影響しない
},
],
});キーポイント
| ポイント | 説明 |
|---|---|
| 選択的適用 | failParentOnFailure: true を持つ子ジョブのみが親の失敗をトリガーする |
| 再帰的動作 | 親にも failParentOnFailure: true が設定されている場合、失敗は祖父母以上に伝播する |
| 即時効果 | 対象の子ジョブが失敗した時点で親ジョブは遅延的に failed 状態に移行する |
動作の仕組み
- grandchild-job-1 が失敗 → child-job-1 が失敗(
failParentOnFailure: true)→ root-job も失敗(child-job-1 もfailParentOnFailure: true) - grandchild-job-2 が失敗 → child-job-1 は影響なし(grandchild-job-2 に
failParentOnFailure未設定) - child-job-2 が失敗 → root-job は影響なし(child-job-2 に
failParentOnFailure未設定)
注意点
- 子ジョブが失敗すると、親ジョブは遅延的に failed 状態に移行します。Worker が親ジョブを処理する際に
UnrecoverableError(メッセージ: "child {childKey} failed")が発生します - このオプションは再帰的に検証されるため、設定に応じて祖父母以上も失敗する可能性があります
- 親ジョブの成功が特定の子ジョブに厳密に依存するワークフローで特に有用です
関連
- ./flows.md — FlowProducer の基本
- ./continue-parent.md — 親ジョブの継続処理
- ./remove-dependency.md — 依存関係の削除
- ./ignore-dependency.md — 依存関係の無視
Flows
BullMQ は親子関係を持つジョブ(フロー)をサポートしており、FlowProducer クラスを使用して任意の深さのツリー構造でジョブを作成・管理できます。親ジョブはすべての子ジョブが正常に完了するまで処理されません。
FlowJob インターフェース
interface FlowJobBase<T> {
name: string;
queueName: string;
data?: any;
prefix?: string;
opts?: Omit<T, 'debounce' | 'deduplication' | 'repeat'>;
children?: FlowChildJob[];
}
type FlowChildJob = FlowJobBase<
Omit<JobsOptions, 'debounce' | 'deduplication' | 'parent' | 'repeat'>
>;
type FlowJob = FlowJobBase<JobsOptions>;基本的なフロー追加
import { FlowProducer } from 'bullmq';
const flowProducer = new FlowProducer();
const flow = await flowProducer.add({
name: 'renovate-interior',
queueName: 'renovate',
children: [
{ name: 'paint', data: { place: 'ceiling' }, queueName: 'steps' },
{ name: 'paint', data: { place: 'walls' }, queueName: 'steps' },
{ name: 'fix', data: { place: 'floor' }, queueName: 'steps' },
],
});上記のコードは4つのジョブをアトミックに追加します。"steps" キューの3つのジョブが完了すると、"renovate" キューの親ジョブが通常のジョブとして処理されます。
子ジョブの結果を取得(getChildrenValues)
子ジョブの Worker が値を返す場合:
import { Worker } from 'bullmq';
const stepsWorker = new Worker('steps', async job => {
await performStep(job.data);
if (job.name === 'paint') {
return 2500;
} else if (job.name === 'fix') {
return 1750;
}
});親 Worker で getChildrenValues メソッドを使って子ジョブの結果を集約できます:
import { Worker } from 'bullmq';
const renovateWorker = new Worker('renovate', async job => {
const childrenValues = await job.getChildrenValues();
const totalCosts = Object.values(childrenValues).reduce(
(prev, cur) => prev + cur,
0,
);
await sendInvoice(totalCosts);
});直列実行(深いツリー構造)
ジョブを直列に実行するには、深いネスト構造を使用します:
import { FlowProducer } from 'bullmq';
const flowProducer = new FlowProducer();
const queueName = 'assembly-line';
const chain = await flowProducer.add({
name: 'car',
data: { step: 'engine' },
queueName,
children: [
{
name: 'car',
data: { step: 'wheels' },
queueName,
children: [{ name: 'car', data: { step: 'chassis' }, queueName }],
},
],
});処理順序: chassis → wheels → engine
Getters
getDependencies
ジョブの直接の依存関係(子ジョブ)を取得します:
const dependencies = await job.getDependencies();特定の種類の子ジョブをページネーション付きで取得:
const { processed, nextProcessedCursor } = await job.getDependencies({
processed: { count: 5, cursor: 0 },
});
const { unprocessed, nextUnprocessedCursor } = await job.getDependencies({
unprocessed: { count: 5, cursor: 0 },
});
const { failed, nextFailedCursor } = await job.getDependencies({
failed: { count: 5, cursor: 0 },
});
const { ignored, nextIgnoredCursor } = await job.getDependencies({
ignored: { count: 5, cursor: 0 },
});getDependenciesCount
子ジョブの種類別カウントを取得:
const { failed, ignored, processed, unprocessed } =
await job.getDependenciesCount();特定の種類のみ取得:
const { failed } = await job.getDependenciesCount({ failed: true });
const { ignored, processed } = await job.getDependenciesCount({
ignored: true,
processed: true,
});getChildrenValues
子ジョブが返した値をすべて取得:
const values = await job.getChildrenValues();parentKey プロパティ
Job クラスに parentKey プロパティがあり、親ジョブの完全修飾キーを保持します。
waiting-children ステート
親ジョブは子ジョブの完了待ちの間 waiting-children ステートになります:
const state = await job.getState();
// state will be "waiting-children"キューオプションの指定
フロー追加時に queueOptions を使ってキューごとのオプションを指定できます:
import { FlowProducer } from 'bullmq';
const flowProducer = new FlowProducer();
const queueName = 'assembly-line';
const chain = await flowProducer.add(
{
name: 'car',
data: { step: 'engine' },
queueName,
children: [
{
name: 'car',
data: { step: 'wheels' },
queueName,
},
],
},
{
queuesOptions: {
[queueName]: {
defaultJobOptions: {
removeOnComplete: true,
},
},
},
},
);ジョブの削除
フロー内のジョブ削除に関する重要なルール:
| ルール | 説明 |
|---|---|
| 親ジョブの削除 | すべての子ジョブも削除される |
| 子ジョブの削除 | 親の依存関係から削除され、最後の子だった場合は親が completed になる |
| 親兼子ジョブの削除 | 上記の両方のルールが適用される |
| ロック中のジョブ | いずれかのジョブがロック中の場合、どのジョブも削除されず例外がスローされる |
await job.remove();
// or
await queue.remove(job.id);注意点
- フローは
FlowProducerクラスを使って追加する必要があります - 親キューと子キューは同じキューである必要はありません
jobIdオプションにコロン:を含めないでください(セパレータとして扱われます)- キューオプションは第2引数で指定してください(インスタンスのコンテキストで定義されるため)
関連
- ./adding-bulks.md — 複数フローの一括追加
- ./get-flow-tree.md — フローツリーの取得
- ./fail-parent.md — 子失敗時の親への影響
- ./continue-parent.md — 親の継続処理
- ./remove-dependency.md — 依存関係の削除
- ./ignore-dependency.md — 依存関係の無視
- ./remove-child-dependency.md — 子依存関係の削除
Get Flow Tree
FlowProducer.getFlow() メソッドを使用して、ジョブとそのすべての子・孫ジョブをツリー構造として取得できます。大規模なフローの可視化やデバッグに有用です。
基本的な使い方
const flow = new FlowProducer({ connection });
const originalTree = await flow.add({
name: 'root-job',
queueName: 'topQueueName',
data: {},
children: [
{
name,
data: { idx: 0, foo: 'bar' },
queueName: 'childrenQueueName',
children: [
{
name,
data: { idx: 4, foo: 'baz' },
queueName: 'grandchildrenQueueName',
},
],
},
{
name,
data: { idx: 2, foo: 'foo' },
queueName: 'childrenQueueName',
},
{
name,
data: { idx: 3, foo: 'bis' },
queueName: 'childrenQueueName',
},
],
});
const { job: topJob } = originalTree;
const tree = await flow.getFlow({
id: topJob.id,
queueName: 'topQueueName',
});
const { children, job } = tree;各 child は job プロパティを持ち、さらに子がある場合は children プロパティも持ちます。
ページネーション(depth / maxChildren)
大量の子ジョブがある場合、取得を制限できます:
const limitedTree = await flow.getFlow({
id: topJob.id,
queueName: 'topQueueName',
depth: 1, // 第1階層の子のみ取得
maxChildren: 2, // ノードごとに最大2つの子を取得
});
const { children, job } = limitedTree;| Name | Type | Description |
|---|---|---|
| id | string | ルートジョブの ID |
| queueName | string | ルートジョブのキュー名 |
| depth | number | 取得するツリーの深さ制限 |
| maxChildren | number | 各ノードの子の最大取得数 |
注意点
- 各
childオブジェクトはjobプロパティを持ち、子がある場合はchildrenプロパティも含まれます depthとmaxChildrenを適切に設定して、大規模フローでのパフォーマンス問題を回避してください
関連
- ./flows.md — FlowProducer の基本
- ./adding-bulks.md — 複数フローの一括追加
Ignore Dependency
ignoreDependencyOnFailure オプションを使用すると、子ジョブが失敗した際に依存関係を無視して親ジョブの処理を続行できます。removeDependencyOnFailure と異なり、失敗情報は保持され、getIgnoredChildrenFailures で取得可能です。
基本的な使い方
const flow = new FlowProducer({ connection });
const originalTree = await flow.add({
name: 'root-job',
queueName: 'topQueueName',
data: {},
children: [
{
name,
data: { idx: 0, foo: 'bar' },
queueName: 'childrenQueueName',
opts: { ignoreDependencyOnFailure: true },
children: [
{
name,
data: { idx: 1, foo: 'bah' },
queueName: 'grandChildrenQueueName',
},
{
name,
data: { idx: 2, foo: 'baz' },
queueName: 'grandChildrenQueueName',
},
],
},
{
name,
data: { idx: 3, foo: 'foo' },
queueName: 'childrenQueueName',
},
],
});失敗情報の取得
const ignoredChildrenFailures =
await originalTree.job.getIgnoredChildrenFailures();注意点
ignoreDependencyOnFailure: trueを持つ子ジョブが失敗すると、親ジョブの依存関係から無視されます- 他に保留中の子ジョブがなければ、親ジョブは waiting 状態に移行します
removeDependencyOnFailureとの違いは、失敗した子の情報がgetIgnoredChildrenFailuresで取得できる点です
関連
- ./flows.md — FlowProducer の基本
- ./remove-dependency.md — 依存関係の削除
- ./remove-child-dependency.md — 子依存関係の手動削除
- ./fail-parent.md — 子失敗時の親の失敗
BullMQ — Flows
| Name | Description | Path |
|---|---|---|
| Flows | FlowProducer による親子ジョブの作成と管理 | ./flows.md |
| Adding flows in bulk | 複数フローのアトミックな一括追加 | ./adding-bulks.md |
| Get Flow Tree | フローツリーの取得と可視化 | ./get-flow-tree.md |
| Fail Parent | 子ジョブ失敗時の親ジョブへの影響 | ./fail-parent.md |
| Continue Parent | 親ジョブの継続処理 | ./continue-parent.md |
| Remove Dependency | 依存関係の削除 | ./remove-dependency.md |
| Ignore Dependency | 依存関係の無視 | ./ignore-dependency.md |
| Remove Child Dependency | 子の依存関係の削除 | ./remove-child-dependency.md |
Remove Child Dependency
removeChildDependency メソッドを使用して、子ジョブから親ジョブへの依存関係を手動で削除できます。対象の子ジョブが最後の保留中の子であった場合、親ジョブは waiting 状態に移行します。
基本的な使い方
const flow = new FlowProducer({ connection });
const originalTree = await flow.add({
name: 'root-job',
queueName: 'topQueueName',
data: {},
children: [
{
name,
data: { idx: 0, foo: 'bar' },
queueName: 'childrenQueueName',
opts: {},
},
],
});
await originalTree.children[0].job.removeChildDependency();注意点
- このメソッドを呼び出すと、既存の親ジョブがあるか検証され、親が存在しない場合はエラーがスローされます
- 対象の子ジョブが最後の保留中の子であった場合、親ジョブは waiting 状態に移行して処理されます
- 既に failed または completed 状態の子ジョブに対しては、unprocessed リストに含まれないため削除は発生しません
関連
- ./flows.md — FlowProducer の基本
- ./remove-dependency.md — 失敗時の自動依存関係削除
- ./ignore-dependency.md — 依存関係の無視
Remove Dependency
removeDependencyOnFailure オプションを使用すると、子ジョブが失敗した際に親ジョブとの依存関係を自動的に削除できます。親ジョブは失敗した子ジョブの完了を待たずに処理を続行します。
基本的な使い方
const flow = new FlowProducer({ connection });
const originalTree = await flow.add({
name: 'root-job',
queueName: 'topQueueName',
data: {},
children: [
{
name,
data: { idx: 0, foo: 'bar' },
queueName: 'childrenQueueName',
opts: { removeDependencyOnFailure: true },
children: [
{
name,
data: { idx: 1, foo: 'bah' },
queueName: 'grandChildrenQueueName',
},
{
name,
data: { idx: 2, foo: 'baz' },
queueName: 'grandChildrenQueueName',
},
],
},
{
name,
data: { idx: 3, foo: 'foo' },
queueName: 'childrenQueueName',
},
],
});注意点
removeDependencyOnFailure: trueを持つ子ジョブが失敗すると、親ジョブの依存関係リストからその子が削除されます- 他に保留中の子ジョブがなければ、親ジョブは waiting 状態に移行して処理されます
- 失敗した子ジョブの結果は
getChildrenValuesには含まれません
関連
- ./flows.md — FlowProducer の基本
- ./fail-parent.md — 子失敗時の親の失敗
- ./ignore-dependency.md — 依存関係の無視
- ./remove-child-dependency.md — 子依存関係の手動削除
Going to Production
BullMQ ベースのアプリケーションを本番環境にデプロイする際の重要な考慮事項とベストプラクティスをまとめています。Redis の永続化設定、メモリポリシー、エラーハンドリング、グレースフルシャットダウンなど、堅牢な運用に不可欠な設定を網羅します。
Redis 永続化(Persistence)
BullMQ は Redis ベースのため、永続化を手動で設定する必要があります。多くのホスティングソリューションではデフォルトで永続化が無効です。
推奨設定: AOF(Append Only File)を有効にする。通常、1 秒ごとの書き込みで十分です。
# redis.conf
appendonly yes
appendfsync everysecMax Memory Policy
Redis をキャッシュとして使用する場合、メモリ上限に達するとキーが削除されますが、BullMQ ではキーの任意削除は許容されません。
必須設定: maxmemory-policy を noeviction に設定する。これが BullMQ の正常動作を保証する唯一の設定です。
# redis.conf
maxmemory-policy noeviction自動再接続
本番環境では Redis との接続が切断される可能性があります。IORedis の以下のオプションを適切に設定することが重要です。
| Option | Queue での推奨値 | Worker での推奨値 | Description |
|---|---|---|---|
retryStrategy | カスタム | カスタム | リトライ間隔の計算関数 |
maxRetriesPerRequest | デフォルト | null | リクエストあたりの最大リトライ回数 |
enableOfflineQueue | false | true(デフォルト) | オフラインキューの有効化 |
retryStrategy
BullMQ のデフォルトでは指数バックオフ(最小 1 秒、最大 20 秒)を使用します。
retryStrategy: function (times: number) {
return Math.max(Math.min(Math.exp(times), 20000), 1000);
}maxRetriesPerRequest
Worker では必ず null に設定します。デフォルトでそう設定されますが、既存の IORedis インスタンスを渡す場合は注意が必要です。
enableOfflineQueue
Queue では false にしてフェイルファストを実現し、Worker ではデフォルト(true)のままにして再接続まで待機させます。
エラーログ
接続問題のデバッグと「unhandled errors」防止のため、error イベントハンドラを設定します。
worker.on('error', (err) => {
logger.error(err, 'Worker error');
});
queue.on('error', (err) => {
logger.error(err, 'Queue error');
});グレースフルシャットダウン
サーバー再起動時にストールジョブを最小限にするため、SIGINT と SIGTERM シグナルをハンドリングします。
const gracefulShutdown = async (signal: string) => {
console.log(`Received ${signal}, closing server...`);
await worker.close();
// その他の非同期クリーンアップ
process.exit(0);
};
process.on('SIGINT', () => gracefulShutdown('SIGINT'));
process.on('SIGTERM', () => gracefulShutdown('SIGTERM'));| Signal | 説明 |
|---|---|
SIGINT | ターミナルで Ctrl+C が押された場合 |
SIGTERM | Kubernetes、PM2 などのオーケストレーションツールからの終了要求 |
ジョブの自動削除
デフォルトでは完了・失敗したジョブは永遠に保持されます。適切な自動削除を設定してください。
const queue = new Queue('myQueue', {
defaultJobOptions: {
removeOnComplete: {
count: 100, // 最新100件の完了ジョブを保持
},
removeOnFail: {
count: 5000, // 最新5000件の失敗ジョブを保持
},
},
});データの保護
ジョブの data フィールドは平文で Redis に保存されます。
- 推奨: 機密データをジョブに含めない
- やむを得ない場合: キューに追加する前にデータを暗号化する
未処理の例外とリジェクション
NodeJS はデフォルトで未処理の例外でクラッシュします。以下のハンドラを設定してください。
process.on('uncaughtException', function (err) {
logger.error(err, 'Uncaught exception');
});
process.on('unhandledRejection', (reason, promise) => {
logger.error({ promise, reason }, 'Unhandled Rejection at: Promise');
});注意点
- AOF 永続化はパフォーマンスに影響するため、ベンチマークで許容範囲を確認すること
maxmemory-policyがnoeviction以外だと、ロックキーが期限前に削除される可能性がある- Worker の
maxRetriesPerRequestをnull以外にすると、Redis コマンドの例外でワーカーが壊れる - グレースフルシャットダウンでもジョブの処理時間がグレース期間を超える場合、ストールは避けられない
- セキュリティは軽視せず、データ漏洩やビジネスへの経済的損害のリスクを真剣に考慮すること
関連
- ./architecture.md
- ./troubleshooting.md
- ./rate-limiting.md
- ./redis-compatibility/redis-compatibility.md
- ./redis-hosting/aws-memorydb.md
BullMQ — Introduction
BullMQ は 4 つのコアクラスを中心に構築されたジョブキューライブラリである。Queue でジョブを登録し、Worker で処理し、QueueEvents でイベントを監視し、FlowProducer で親子関係のあるワークフローを構築する。
コアクラス
Queue
キューを表現するクラス。ジョブの追加、一時停止、クリーニング、データ取得など基本的な操作を提供する。
import { Queue } from 'bullmq';
const myQueue = new Queue('Paint');
await myQueue.add('cars', { color: 'blue' });Worker
キューからジョブを取り出して処理するインスタンス。複数の Worker を異なる Node.js プロセスやマシンで同時に実行でき、ジョブを completed または failed としてマークする。
import { Worker } from 'bullmq';
const worker = new Worker('Paint', async job => {
// process job
console.log(job.data);
});QueueEvents
キューに関連するイベントの監視とハンドリングを可能にするクラス。ジョブの完了、失敗、進捗などのイベントをリッスンできる。
import { QueueEvents } from 'bullmq';
const queueEvents = new QueueEvents('Paint');
queueEvents.on('completed', ({ jobId }) => {
console.log(`Job ${jobId} completed`);
});FlowProducer
複雑なジョブワークフローのオーケストレーションを実現するクラス。親子関係のあるジョブの依存関係を定義できる。
import { FlowProducer } from 'bullmq';
const flowProducer = new FlowProducer();
const flow = await flowProducer.add({
name: 'renovate-interior',
queueName: 'renovate',
children: [
{ name: 'paint', data: { place: 'ceiling' }, queueName: 'steps' },
{ name: 'paint', data: { place: 'walls' }, queueName: 'steps' },
],
});ジョブのライフサイクル
ジョブはユーザーが定義したデータ構造であり、キューに保存される。基本的な流れは以下の通り:
1. 追加: Queue.add() でジョブをキューに投入 2. 待機: Redis 内のリストでワーカーを待つ(waiting 状態) 3. 処理: Worker がジョブを取得し処理する(active 状態) 4. 完了/失敗: 処理結果に応じて completed または failed に遷移
Worker がジョブ追加時に稼働していなくても、接続された時点でキュー内のジョブを自動的に処理する。
注意点
- BullMQ は分散処理をサポートしており、複数マシンにまたがるスケーラブルなアーキテクチャを構築可能
- ジョブデータは Redis に保存されるため、シリアライズ可能なデータである必要がある
- Queue, Worker, QueueEvents はそれぞれ独立した Redis 接続を使用できる
関連
- Connections
- Queues
- Workers
BullMQ — Job Schedulers
Job Scheduler は指定された repeat 設定に基づいてジョブを生成するファクトリとして機能します。固定間隔、cron 式、カスタム要件など、さまざまなシナリオに対応できます。歴史的な理由から、Job Scheduler が生成するジョブは「Repeatable Jobs」と呼ばれることがあります。
基本的な使い方
upsertJobScheduler メソッドでスケジューラを作成します。
import { Queue } from 'bullmq';
const queue = new Queue('Paint');
// 1秒ごとにジョブを生成するスケジューラを作成
const firstJob = await queue.upsertJobScheduler('my-scheduler-id', {
every: 1000,
});この例では、1秒ごとに新しいジョブを生成する Job Scheduler が作成されます。最初のジョブは「delayed」状態で返され、1秒後に処理されます。
Job テンプレートの使用
ジョブに標準的な name、data、options を定義するテンプレートを指定できます。
// 毎日 3:15(午前)にジョブを作成
const firstJob = await queue.upsertJobScheduler(
'my-scheduler-id',
{ pattern: '0 15 3 * * *' },
{
name: 'my-job-name',
data: { foo: 'bar' },
opts: {
backoff: 3,
attempts: 5,
removeOnFail: 1000,
},
},
);同じスケジューラ ID で upsertJobScheduler を再度呼び出すことで、repeat オプションやジョブテンプレートの設定を更新できます。
upsertJobScheduler パラメータ
| Name | Type | Description |
|---|---|---|
jobSchedulerId | string | スケジューラの一意な識別子 |
repeatOpts | object | 繰り返し設定(every, pattern など) |
template | object | ジョブテンプレート(name, data, opts) |
注意点
- Upsert vs Add: 本番デプロイメントでの管理を簡素化するために
upsertを使用。スケジューラの重複なく更新または作成が保証される - ジョブ生成レート: スケジューラは最後のジョブが処理を開始したときにのみ新しいジョブを生成する。キューが混雑している場合やワーカー/コンカレンシーが不足している場合、指定間隔より低頻度になる可能性がある
- ジョブの状態: Job Scheduler がジョブを生成している限り、常に1つのジョブが「Delayed」状態でスケジューラに関連付けられている
- Job Scheduler が生成するジョブには特別なジョブ ID が付与されるため、カスタムジョブ ID は指定できない。ジョブの識別にはジョブ名を使用すること
関連
- Repeat Strategies
- Repeat Options
- Manage Job Schedulers
- ../jobs/repeatable.md
BullMQ — Manage Job Schedulers
Job Scheduler のライフサイクル管理は、効率的なバックグラウンドタスクの維持に不可欠です。upsertJobScheduler に加えて、removeJobScheduler と getJobSchedulers メソッドでスケジューラの削除と一覧取得が行えます。
removeJobScheduler(スケジューラの削除)
不要になったスケジューラや、非アクティブ/廃止されたスケジューラを削除してリソースを最適化します。
import { Queue } from 'bullmq';
const queue = new Queue('Paint');
// スケジューラ 'scheduler-123' を削除
const result = await queue.removeJobScheduler('scheduler-123');
console.log(
result ? 'Scheduler removed successfully' : 'Missing Job Scheduler',
);| 戻り値 | Description |
|---|---|
true | 指定 ID の Job Scheduler が存在し、削除された |
false | 指定 ID の Job Scheduler が存在しなかった |
getJobSchedulers(スケジューラ一覧の取得)
指定した範囲内のすべてのスケジューラ設定を取得します。モニタリングやダッシュボードでの使用に最適です。
// 最初の10件のスケジューラを次回実行時刻の昇順で取得
const schedulers = await queue.getJobSchedulers(0, 9, true);
console.log('Current job schedulers:', schedulers);getJobSchedulers パラメータ
| Name | Type | Description |
|---|---|---|
start | number | 取得開始位置 |
end | number | 取得終了位置 |
asc | boolean | true で次回実行時刻の昇順 |
getJobScheduler(個別スケジューラの取得)
ID を指定して特定のスケジューラの設定を取得します。
const scheduler = await queue.getJobScheduler('test');
console.log('Current job scheduler:', scheduler);注意点
removeJobSchedulerはスケジューラが存在しない場合falseを返す(エラーにはならない)getJobSchedulersはページネーションに対応しており、大量のスケジューラがある場合でも効率的に取得可能- スケジューラの更新には
upsertJobSchedulerを使用する(同じ ID で呼び出すと既存設定が上書きされる) - スケジューラのレポートやダッシュボード生成には
getJobSchedulersが有用
関連
- Job Schedulers
- Repeat Strategies
- Repeat Options
BullMQ — Job Schedulers
| Name | Description | Path |
|---|---|---|
| Job Schedulers | JobScheduler の概要と基本的な使い方 | ./job-schedulers.md |
| Repeat Strategies | cron / every による繰り返し戦略 | ./repeat-strategies.md |
| Repeat Options | 繰り返しオプションの詳細設定 | ./repeat-options.md |
| Manage Job Schedulers | ジョブスケジューラの管理(取得・削除) | ./manage-job-schedulers.md |
BullMQ — Repeat Options
すべての Job Scheduler で使用可能な繰り返しオプションについて解説します。これらのオプションで繰り返しの開始・終了日時、回数制限、即時実行などを制御できます。
オプション一覧
| Name | Type | Description |
|---|---|---|
every | number | 繰り返し間隔(ミリ秒) |
pattern | string | cron 式 |
startDate | Date | スケジュール開始日。この日以降からジョブが生成される |
endDate | Date | スケジュール終了日。この日以降はジョブが生成されない |
limit | number | 最大繰り返し回数。この回数に達すると以降のジョブは生成されない |
immediately | boolean | 即時実行(非推奨: v5.19.0 以降はデフォルトで即時実行) |
Start Date(開始日)
将来の日付を指定し、その日からスケジュールを開始します。
const { Queue } = require('bullmq');
const connection = { host: 'localhost', port: 6379 };
const myQueue = new Queue('my-dated-jobs', { connection });
await myQueue.upsertJobScheduler(
'start-later-job',
{
every: 60000, // 毎分
startDate: new Date('2024-10-15T00:00:00Z'),
},
{
name: 'timed-start-job',
data: { message: 'Starting later' },
},
);End Date(終了日)
ジョブの繰り返しを終了する日を指定します。
await myQueue.upsertJobScheduler(
'end-soon-job',
{
every: 60000,
endDate: new Date('2024-11-01T00:00:00Z'),
},
{
name: 'timed-end-job',
data: { message: 'Ending soon' },
},
);Limit(回数制限)
繰り返し回数を制限します。count がこの制限に達すると、それ以上のジョブは生成されません。
await myQueue.upsertJobScheduler(
'limited-job',
{
every: 10000, // 10秒ごと
limit: 10, // 最大10回
},
{
name: 'limited-execution-job',
data: { message: 'Limited runs' },
},
);Immediately(即時実行)
スケジュールに関係なく、追加直後にジョブを実行します。
await myQueue.upsertJobScheduler(
'immediate-job',
{
every: 86400000, // 1日1回
immediately: true,
},
{
name: 'instant-job',
data: { message: 'Immediate start' },
},
);Note:everyオプションは固定時間間隔に基づいてスケジュールされます。例えば 2000ms の間隔を設定すると、ジョブはクロックの偶数秒(0, 2, 4, 6, 8...)にトリガーされます。immediatelyを使うと、間隔のアラインメントに関係なく最初のジョブが即座に実行されます。
注意点
- バージョン 5.19.0 以降、
immediatelyオプションは非推奨。新規挿入された Job Scheduler の最初の繰り返しは常に即時実行される(既存スケジューラのupsertJobSchedulerはevery間隔に従う) startDateとendDateを組み合わせることで、特定期間のみ有効なスケジュールを設定できるlimitに達した後は、そのスケジューラからジョブは生成されなくなる- 月次など長い間隔を設定した場合、
immediatelyなしでは月初まで待機する必要があった(v5.19.0 以前)
関連
- Job Schedulers
- Repeat Strategies
- Manage Job Schedulers
BullMQ — Repeat Strategies
BullMQ には Repeatable ジョブを作成するための2つの定義済み戦略(every と cron)があり、さらにカスタム戦略を定義することも可能です。
"Every" 戦略
指定したミリ秒間隔でジョブを繰り返し生成するシンプルな戦略です。
const { Queue, Worker } = require('bullmq');
const connection = { host: 'localhost', port: 6379 };
const myQueue = new Queue('my-repeatable-jobs', { connection });
// 10秒ごとにジョブを繰り返し
await myQueue.upsertJobScheduler(
'repeat-every-10s',
{
every: 10000,
},
{
name: 'every-job',
data: { jobData: 'data' },
opts: {},
},
);
const worker = new Worker(
'my-repeatable-jobs',
async job => {
console.log(`Processing job ${job.id} with data: ${job.data.jobData}`);
},
{ connection },
);"Cron" 戦略
cron-parser ライブラリを使用した cron 式によるスケジューリングです。自動レポートやメンテナンスタスクなど、正確な時刻での実行に最適です。
cron 式のフォーマット
* * * * * *
| | | | | |
| | | | | +-- day of week (0-7, 1L-7L, 0 or 7 = Sunday)
| | | | +------- month (1-12)
| | | +------------ day of month (1-31, L = last day)
| | +----------------- hour (0-23)
| +---------------------- minute (0-59)
+--------------------------- second (0-59, optional)const { Queue, Worker } = require('bullmq');
const connection = { host: 'localhost', port: 6379 };
const myQueue = new Queue('my-cron-jobs', { connection });
// 月~金 毎朝 9:00 に実行
await myQueue.upsertJobScheduler(
'weekday-morning-job',
{
pattern: '0 0 9 * * 1-5',
},
{
name: 'cron-job',
data: { jobData: 'morning data' },
opts: {},
},
);
const worker = new Worker(
'my-cron-jobs',
async job => {
console.log(
`Processing job ${job.id} at ${new Date()} with data: ${job.data.jobData}`,
);
},
{ connection },
);カスタム戦略
独自のスケジューリングロジックを定義できます。repeat strategy は pattern と最新ジョブのミリ秒に基づいて次のタイムスタンプを返します。1つのキューにつき1つの repeatStrategy のみ定義可能です。
RRULE を使用したカスタム戦略の例:
import { Queue, Worker } from 'bullmq';
import { rrulestr } from 'rrule';
const settings = {
repeatStrategy: (millis: number, opts: RepeatOptions, _jobName: string) => {
const currentDate =
opts.startDate && new Date(opts.startDate) > new Date(millis)
? new Date(opts.startDate)
: new Date(millis);
const rrule = rrulestr(opts.pattern);
if (rrule.origOptions.count && !rrule.origOptions.dtstart) {
throw new Error('DTSTART must be defined to use COUNT with rrule');
}
const next_occurrence = rrule.after(currentDate, false);
return next_occurrence?.getTime();
},
};
const myQueue = new Queue('Paint', { settings });
await myQueue.upsertJobScheduler(
'collibris',
{
pattern: 'RRULE:FREQ=SECONDLY;INTERVAL=10;WKST=MO',
},
{
data: { color: 'green' },
},
);
await myQueue.upsertJobScheduler(
'pingeons',
{
pattern: 'RRULE:FREQ=SECONDLY;INTERVAL=20;WKST=MO',
},
{
data: { color: 'gray' },
},
);
const worker = new Worker(
'Paint',
async () => {
doSomething();
},
{ settings },
);注意点
- cron 式はタイムゾーンの違いや夏時間の切り替えをシームレスに処理できる
- cron 式のオプションの秒フィールドにより、標準 cron より精密なスケジューリングが可能
repeatStrategy設定は Queue と Worker の両方に提供する必要がある(初回はキューで次回タイミングを計算し、以降はワーカーが引き継ぐため)- 1つのキューにつき1つの
repeatStrategyのみ定義可能
関連
- Job Schedulers
- Repeat Options
- Manage Job Schedulers
BullMQ — Deduplication
BullMQ のデデュプリケーション(重複排除)は、特定の識別子に基づいてジョブの実行を制御し、指定期間内またはジョブ完了/失敗まで同じ識別子の新しいジョブがキューに追加されないようにする機能です。重複追加の試行時には deduplicated イベントが発火します。
Simple モード
Simple モードでは、ジョブの完了または失敗まで重複排除が継続します。ジョブが未完了の間、同じ deduplication ID を持つ後続のジョブは無視されます。
// ジョブが完了/失敗するまで重複排除される
await myQueue.add(
'house',
{ color: 'white' },
{ deduplication: { id: 'customValue' } },
);Throttle モード
Throttle モードでは、TTL(Time to Live)を指定して一定期間内の重複を防止します。
import { Queue } from 'bullmq';
const myQueue = new Queue('Paint');
// 5秒間の重複排除
await myQueue.add(
'house',
{ color: 'white' },
{ deduplication: { id: 'customValue', ttl: 5000 } },
);Debounce モード
Debounce モードは delay と ttl を組み合わせ、extend と replace オプションを true にすることで実現します。同じ deduplication ID のジョブが追加されると、前のジョブを新しいジョブで置き換え、TTL もリセットされます。
import { Queue, Worker } from 'bullmq';
const myQueue = new Queue('Paint');
const worker = new Worker('Paint', async () => {});
worker.once('completed', job => {
console.log(job.data.color); // `white 10`
});
// Debounce モードで10件のジョブを追加
for (let i = 1; i < 11; i++) {
await myQueue.add(
'house1',
{ color: `white ${i}` },
{
deduplication: {
id: 'customValue',
ttl: 5000,
extend: true,
replace: true,
},
delay: 5000,
},
);
}deduplication オプション
| Name | Type | Description |
|---|---|---|
id | string | 重複排除の識別子。ジョブデータのハッシュやサブセットで生成 |
ttl | number | Throttle/Debounce モードでの有効期間(ミリ秒) |
extend | boolean | TTL を延長するかどうか(Debounce モード用) |
replace | boolean | 既存ジョブを新しいジョブで置き換えるか(Debounce モード用) |
Deduplicated イベント
重複排除が発生した際に QueueEvents クラスから deduplicated イベントを受信できます。
import { QueueEvents } from 'bullmq';
const queueEvents = new QueueEvents('myQueue');
queueEvents.on(
'deduplicated',
({ jobId, deduplicationId, deduplicatedJobId }, id) => {
console.log(
`Job ${deduplicatedJobId} was deduplicated due to existing job ${jobId} with deduplication ID ${deduplicationId}`,
);
},
);Deduplication Job ID の取得
const jobId = await myQueue.getDeduplicationJobId('customValue');Deduplication キーの削除
TTL 終了前やジョブ完了前に重複排除を停止したい場合は、キーを手動で削除できます。
// キューから削除
await myQueue.removeDeduplicationKey('customValue');
// 特定のジョブから削除
const isRemoved = await job.removeDeduplicationKey();注意点
- 手動でジョブを削除(
job.remove()の呼び出し)すると重複排除が無効になる - deduplication ID はジョブを一意に表す識別子として設計すること(ジョブデータ全体またはサブセットのハッシュなど)
- Simple モードは長時間実行ジョブや重複実行を防ぎたい重要な処理に適している
- Throttle モードは短時間に大量の同一リクエストが発生するシナリオに適している
- Debounce モードは最新のデータのみを処理したい場合に適している
関連
- Job IDs
- Delayed
- Jobs
BullMQ — Delayed Jobs
遅延ジョブは即座に処理されず、特別な「delayed set」に入り、指定された遅延時間が経過した後に通常のジョブとして処理されます。
基本的な使い方
delay オプションにミリ秒を指定してジョブを追加します。
import { Queue } from 'bullmq';
const myQueue = new Queue('Paint');
// 5秒後に処理されるジョブ
await myQueue.add('house', { color: 'white' }, { delay: 5000 });特定の日時にスケジュール
ターゲット時刻までのミリ秒数を計算して delay に渡します。
const targetTime = new Date('03-07-2035 10:30');
const delay = Number(targetTime) - Number(new Date());
await myQueue.add('house', { color: 'white' }, { delay });遅延の変更(changeDelay)
ジョブ作成後に changeDelay() メソッドで遅延時間を変更できます。新しい遅延は現在時刻から再計算されます。
const job = await Job.create(queue, 'test', { foo: 'bar' }, { delay: 2000 });
await job.changeDelay(4000);オプション
| Name | Type | Description |
|---|---|---|
delay | number | 処理開始までの遅延時間(ミリ秒) |
注意点
- 指定した遅延時間ちょうどにジョブが処理される保証はない(ワーカーの可用性や同時遅延ジョブの数に依存する)
changeDelay()は delayed 状態のジョブにのみ適用可能。他の状態のジョブに対して呼び出すとエラーになる- BullMQ 2.0 以降、
QueueSchedulerは不要
関連
- Repeatable
- FIFO
- Jobs
FIFO
FIFO(First-In, First-Out)は BullMQ のデフォルトのジョブ処理順序。ジョブは追加された順番に処理される。
基本的な使い方
import { Queue } from 'bullmq';
const queue = new Queue('my-queue');
// FIFO はデフォルト動作のため、特別なオプションは不要
await queue.add('job1', { data: 'first' });
await queue.add('job2', { data: 'second' });
await queue.add('job3', { data: 'third' });
// 処理順序: job1 → job2 → job3defaultJobOptions での設定
const queue = new Queue('my-queue', {
defaultJobOptions: {
removeOnComplete: true,
},
});注意点
- FIFO はデフォルト動作のため、明示的な設定は不要
- 遅延ジョブや優先度付きジョブを追加すると、厳密な FIFO 順序は崩れる
- 複数ワーカーが並行処理する場合、個々のジョブの完了順序は FIFO とは限らない
関連
- LIFO
- Prioritized
- Jobs
BullMQ — Getters
ジョブはライフサイクルの中でさまざまなステータスを経由します。BullMQ は各ステータスからジョブ情報を取得するためのメソッドを提供しています。
Job Counts(ジョブ数の取得)
指定したステータスのジョブ数を取得します。
import { Queue } from 'bullmq';
const myQueue = new Queue('Paint');
const counts = await myQueue.getJobCounts('wait', 'completed', 'failed');
// { wait: number, completed: number, failed: number }利用可能なステータス
| Status | Description |
|---|---|
completed | 完了したジョブ |
failed | 失敗したジョブ |
delayed | 遅延中のジョブ |
active | 処理中のジョブ |
wait | 待機中のジョブ |
waiting-children | 子ジョブの完了を待っているジョブ |
prioritized | 優先度付きジョブ |
paused | 一時停止中のジョブ |
repeat | 繰り返しジョブ |
Get Jobs(ジョブの取得)
ページネーション形式でジョブを取得します。
// 最も古い100件の completed ジョブを取得
const completed = await myQueue.getJobs(['completed'], 0, 100, true);getJobs パラメータ
| Name | Type | Description |
|---|---|---|
types | string[] | 取得するジョブのステータス配列 |
start | number | 取得開始位置(0始まり) |
end | number | 取得終了位置 |
asc | boolean | true で古い順、false で新しい順 |
注意点
getJobCountsは複数のステータスを同時に指定できるgetJobsはページネーションに対応しており、大量のジョブがある場合でも効率的に取得可能- Repeatable Job の設定情報は
getJobs()には表示されない。getRepeatableJobs()を使用すること
関連
- Jobs
- Prioritized
- Removing Jobs
BullMQ — Job Data
すべてのジョブはカスタムデータを持つことができます。データはジョブの data 属性に格納され、ジョブ追加時に第2引数として渡します。
基本的な使い方
import { Queue } from 'bullmq';
const myQueue = new Queue('paint');
const job = await myQueue.add('wall', { color: 'red' });
job.data; // { color: 'red' }データの更新
ジョブ追加後にデータを変更する場合は updateData メソッドを使用します。
const job = await Job.create(queue, 'wall', { color: 'red' });
await job.updateData({
color: 'blue',
});
job.data; // { color: 'blue' }注意点
- ジョブデータは JSON シリアライズ可能なオブジェクトである必要がある
- データサイズが大きいと Redis のメモリ使用量に影響するため、必要最小限のデータを格納すること
- 大きなペイロード(ファイル内容など)は外部ストレージに保存し、ジョブデータには参照(URL や ID)のみを含めるのがベストプラクティス
updateDataメソッドでジョブデータを完全に置き換えることが可能
関連
- Jobs
- FIFO
- Getters
Job Ids
BullMQ はデフォルトで一意なジョブ ID を自動生成するが、jobId オプションでカスタム ID を指定することも可能。カスタム ID は重複排除やべき等性の実装に利用できる。
基本的な使い方
import { Queue } from 'bullmq';
const queue = new Queue('my-queue');
// カスタム ID を指定
await queue.add('my-job', { foo: 'bar' }, { jobId: 'custom-id-123' });重複排除
同じ jobId のジョブが既にキューに存在する場合、新しいジョブは無視され duplicated イベントが発行される。
await queue.add('my-job', { data: 1 }, { jobId: 'unique-1' });
await queue.add('my-job', { data: 2 }, { jobId: 'unique-1' }); // 無視される
// QueueEvents で重複を検知
const queueEvents = new QueueEvents('my-queue');
queueEvents.on('duplicated', ({ jobId }) => {
console.log(`Job ${jobId} was duplicated`);
});注意点
- ジョブ ID にコロン
:を含めることはできない(内部的にセパレータとして使用) removeOnComplete/removeOnFailで削除されたジョブの ID は再利用可能- カスタム ID を使う場合、一意性の管理はアプリケーション側の責任
- Flow 内のジョブに
jobIdを指定する場合もコロン制約に注意
関連
- Deduplication
- Auto-removal of jobs
- Jobs
Jobs
キューは異なる種類のジョブを保持でき、それぞれの種類によって処理方法やタイミングが異なる。同一キュー内で FIFO、LIFO、遅延、優先度付きなど異なるタイプのジョブを混在させることが可能。
基本的な使い方
import { Queue } from 'bullmq';
const queue = new Queue('my-queue');
// 基本的なジョブの追加
await queue.add('job-name', { foo: 'bar' });
// オプション付きジョブの追加
await queue.add('job-name', { foo: 'bar' }, {
delay: 5000, // 5秒遅延
priority: 1, // 優先度
attempts: 3, // リトライ回数
backoff: {
type: 'exponential',
delay: 1000,
},
});主なジョブオプション
| Name | Type | Description |
|---|---|---|
delay | number | ジョブの実行を遅延させるミリ秒数 |
priority | number | 優先度(小さい値ほど優先) |
lifo | boolean | LIFO 順序で処理するか |
attempts | number | 失敗時のリトライ回数 |
backoff | object | リトライ時のバックオフ設定 |
removeOnComplete | `boolean \ | number \ |
removeOnFail | `boolean \ | number \ |
jobId | string | カスタムジョブ ID |
timestamp | number | ジョブのタイムスタンプ |
parent | object | 親ジョブの指定(Flow で使用) |
注意点
- 同一キューで異なる種類のジョブを自由に混在可能
- ジョブデータは JSON シリアライズ可能である必要がある
- ジョブ名はワーカー内で処理を分岐する際に利用できる
関連
- FIFO
- LIFO
- Delayed
- Prioritized
- Queues
- Workers
LIFO
LIFO(Last-In, First-Out)はジョブを追加された逆順で処理する方式。lifo: true オプションで有効化する。
基本的な使い方
import { Queue } from 'bullmq';
const queue = new Queue('my-queue');
await queue.add('job1', { data: 'first' });
await queue.add('job2', { data: 'second' }, { lifo: true });
// job2 が先に処理される(最後に追加されたジョブが最優先)defaultJobOptions での設定
const queue = new Queue('my-queue', {
defaultJobOptions: {
lifo: true,
},
});
// このキューに追加される全ジョブが LIFO で処理される
await queue.add('job1', { data: 'first' });
await queue.add('job2', { data: 'second' });
// 処理順序: job2 → job1注意点
- LIFO ジョブは待機リストの先頭に挿入される
- FIFO ジョブと LIFO ジョブを同一キューで混在可能
- 優先度付きジョブと LIFO を同時に使用する場合、優先度が先に評価される
- スタック的な処理パターン(最新データの優先処理)に有用
関連
- FIFO
- Prioritized
- Jobs
BullMQ — Prioritized Jobs
ジョブに priority オプションを指定すると、FIFO や LIFO のパターンではなく優先度に基づいて処理順序が決定されます。優先度値が小さいほど高い優先度を示します(1 が最高優先度)。
基本的な使い方
import { Queue } from 'bullmq';
const myQueue = new Queue('Paint');
await myQueue.add('wall', { color: 'pink' }, { priority: 10 });
await myQueue.add('wall', { color: 'brown' }, { priority: 5 });
await myQueue.add('wall', { color: 'blue' }, { priority: 7 });
// 処理順: brown (5) → blue (7) → pink (10)オプション
| Name | Type | Description |
|---|---|---|
priority | number | 優先度レベル(1 ~ 2,097,152)。小さいほど高優先度 |
優先度の変更(changePriority)
ジョブ追加後に優先度を変更できます。
const job = await Job.create(queue, 'test2', { foo: 'bar' }, { priority: 16 });
// 優先度を 1 に変更
await job.changePriority({
priority: 1,
});LIFO オプションと組み合わせることも可能です。
const job = await Job.create(queue, 'test2', { foo: 'bar' }, { priority: 16 });
await job.changePriority({
lifo: true,
});Prioritized ジョブの取得
const jobs = await queue.getJobs(['prioritized']);
// または専用メソッド
const jobs2 = await queue.getPrioritized();優先度別カウントの取得
const counts = await queue.getCountsPerPriority([1, 0]);
/*
{
'1': 11, // priority 1 のジョブ数
'0': 10 // priority 0(waiting)のジョブ数
}
*/注意点
- 優先度付きジョブの追加は他のジョブタイプより遅い操作で、計算量は
O(log(n))(n は prioritized set のジョブ数) - 優先度の範囲は
1~2,097,152。値が小さいほど高い優先度 - priority が未指定のジョブは最高優先度として扱われ、priority が指定されたジョブより先に処理される
- 同じ優先度のジョブは FIFO 順で処理される
関連
- FIFO
- LIFO
- Getters
BullMQ — Jobs
| Name | Description | Path |
|---|---|---|
| Jobs | ジョブの概要、種類の混在 | ./jobs.md |
| FIFO | 先入れ先出し順序のジョブ処理 | ./fifo.md |
| LIFO | 後入れ先出し順序のジョブ処理 | ./lifo.md |
| Job Ids | カスタムジョブ ID の指定 | ./job-ids.md |
| Job Data | ジョブデータの設定とサイズの考慮 | ./job-data.md |
| Deduplication | ジョブの重複排除 | ./deduplication.md |
| Delayed | 遅延実行ジョブ | ./delayed.md |
| Repeatable | 繰り返しジョブ(cron / every) | ./repeatable.md |
| Prioritized | 優先度付きジョブ | ./prioritized.md |
| Removing Jobs | ジョブの削除 | ./removing-jobs.md |
| Retrying Jobs | 失敗ジョブのリトライ | ./retrying-jobs.md |
| Stalled | 停滞ジョブ | ./stalled.md |
| Getters | ジョブデータの取得メソッド | ./getters.md |
BullMQ — Removing Jobs
不正なデータを持つジョブの削除など、キューからジョブを手動で削除する必要がある場合に使用します。
基本的な使い方
import { Queue } from 'bullmq';
const queue = new Queue('paint');
const job = await queue.add('wall', { color: 1 });
await job.remove();親ジョブがある場合
子ジョブを削除する際、親ジョブへの影響は2つのケースがあります。
| ケース | 親の状態 | 動作 |
|---|---|---|
| 保留中の依存関係なし | waiting に移動 | 親ジョブの処理が試行される |
| 保留中の依存関係あり | waiting-children のまま | 親ジョブは他の子の完了を待つ |
Note: 削除時に子ジョブが completed 状態だった場合、処理済みの値は親の processed hset に保持されます。
保留中の依存関係がある場合
保留中のすべての子孫ジョブを先に削除しようと試みます。
Warning: 子ジョブのいずれかがロックされている場合、削除プロセスは停止します。
注意点
- ロックされたジョブ(active 状態)は削除できない。削除を試みるとエラーがスローされる
- 親ジョブとの依存関係がある場合、子ジョブの削除は親の状態に影響する
- 子孫ジョブにロックされているものがある場合、削除処理全体が中断される
関連
- Jobs
- Stalled
- Getters
BullMQ — Repeatable Jobs
Repeatable ジョブは、一度キューに追加するだけで事前定義されたスケジュールに従い繰り返し実行される特別なメタジョブです。cron 式またはミリ秒間隔で繰り返しパターンを指定できます。
Note: BullMQ バージョン 5.16.0 以降では、これらの API は非推奨となり、より堅牢な Job Schedulers が推奨されています。
基本的な使い方
repeat オプションを指定してジョブを追加すると、Repeatable Job 設定と最初の遅延ジョブが即座に作成されます。
import { Queue } from 'bullmq';
const myQueue = new Queue('Paint');
// 毎日 3:15(午前)に実行
await myQueue.add(
'submarine',
{ color: 'yellow' },
{
repeat: {
pattern: '0 15 3 * * *',
},
},
);
// 10秒ごとに実行(最大100回)
await myQueue.add(
'bird',
{ color: 'bird' },
{
repeat: {
every: 10000,
limit: 100,
},
},
);repeat オプション
| Name | Type | Description |
|---|---|---|
pattern | string | cron 式(cron-parser の "unix cron w/ optional seconds" 形式) |
every | number | 繰り返し間隔(ミリ秒) |
limit | number | 最大繰り返し回数 |
key | string | カスタム Repeatable キー(同じ repeat オプションのジョブを区別する) |
Repeatable ジョブの削除
import { Queue } from 'bullmq';
const repeat = { pattern: '*/1 * * * * *' };
const myQueue = new Queue('Paint');
const job1 = await myQueue.add('red', { foo: 'bar' }, { repeat });
const job2 = await myQueue.add('blue', { foo: 'baz' }, { repeat });
// repeatJobKey で削除
const isRemoved1 = await myQueue.removeRepeatableByKey(job1.repeatJobKey);
// ジョブ名と repeat オプションで削除
const isRemoved2 = await myQueue.removeRepeatable('blue', repeat);Repeatable ジョブの一覧取得
const repeatableJobs = await myQueue.getRepeatableJobs();カスタム Repeatable キー
同じ repeat オプションを持つジョブを区別するためにカスタムキーを使用できます。
import { Queue } from 'bullmq';
const myQueue = new Queue('Paint', { connection });
await myQueue.add(
'bird',
{ color: 'gray' },
{
repeat: {
every: 10_000,
key: 'colibri',
},
},
);
await myQueue.add(
'bird',
{ color: 'brown' },
{
repeat: {
every: 10_000,
key: 'eagle',
},
},
);カスタムキーによる既存ジョブの更新
同じキーで新しい repeatable ジョブを追加すると、既存の設定が更新されます。
// 間隔を10秒から25秒に変更
await myQueue.add(
'bird',
{ color: 'turquoise' },
{
repeat: {
every: 25_000,
key: 'eagle',
},
},
);カスタム Repeat Strategy
デフォルトの cron-parser ベースの戦略を変更し、独自のスケジューリングロジック(例: RRULE)を定義できます。
import { Queue, Worker } from 'bullmq';
import { rrulestr } from 'rrule';
const settings = {
repeatStrategy: (millis, opts) => {
const currentDate =
opts.startDate && new Date(opts.startDate) > new Date(millis)
? new Date(opts.startDate)
: new Date(millis);
const rrule = rrulestr(opts.pattern);
if (rrule.origOptions.count && !rrule.origOptions.dtstart) {
throw new Error('DTSTART must be defined to use COUNT with rrule');
}
const next_occurrence = rrule.after(currentDate, false);
return next_occurrence?.getTime();
},
};
const myQueue = new Queue('Paint', { settings });
const worker = new Worker(
'Paint',
async () => {
doSomething();
},
{ settings },
);遅い Repeatable ジョブ
繰り返し頻度よりもジョブの処理時間が長い場合、ワーカー数が不足すると期待通りの頻度で処理されないことがあります。例えば 1 秒ごとのジョブで処理に 5 秒かかる場合、5 台のワーカーが必要です。
注意点
- BullMQ は同じ repeat オプションの重複した repeatable ジョブを追加しない
- ワーカーが稼働していない間、repeatable ジョブは蓄積されない
- Repeatable Job 設定はジョブではないため
getJobs()には表示されない。管理にはgetRepeatableJobs()を使用する repeatStrategy設定は Queue と Worker の両方に提供する必要がある(初回はキューで次回タイミングを計算し、以降はワーカーが引き継ぐため)- repeat strategy 関数はオプションの第3引数として
jobNameを受け取る - repeatable ジョブの
jobIdは一意 ID の生成に使用される(通常のジョブとは異なる動作)
関連
- ../job-schedulers/job-schedulers.md
- Delayed
- FIFO
BullMQ — Retrying Jobs
retry メソッドは、completed または failed 状態のジョブを手動で再処理するための機能です。自動リトライメカニズムとは異なり、明示的な手動介入で使用します。
基本的な使い方
import { Queue, Job } from 'bullmq';
const queue = new Queue('paint');
// 失敗したジョブをリトライ
const job = await queue.getJob('job-id');
await job.retry('failed');リトライオプション
| Name | Type | Description |
|---|---|---|
resetAttemptsMade | boolean | ジョブのリトライ許容回数をリセットする |
resetAttemptsStarted | boolean | active 状態への遷移カウンターをリセットする |
主なユースケース
- 外部の一時的な障害が解消された後の手動介入
- 同一データで完了したジョブを再処理
- システム障害後のワークフロー復旧
エラーコード
| コード | Description |
|---|---|
-1 | ジョブが存在しない |
-3 | ジョブが期待される状態にない |
注意点
retryはcompletedまたはfailed状態のジョブにのみ使用可能- リトライ実行時、ジョブは waiting キューに戻され、失敗理由や処理タイムスタンプなどの関連プロパティがクリアされる
attemptsMadeをリセットせずにリトライし、既にリトライ回数を使い切っている場合、ジョブは再処理時に即座に失敗する- TypeScript、Python、Elixir の各実装で利用可能
関連
- Stalled
- Removing Jobs
- Jobs