BullMQとは:Node.js(TypeScript)のジョブキューライブラリ入門
BullMQとは
BullMQはRedisをバックエンドに使うNode.js向けのジョブキューライブラリです。メール送信・画像変換・外部API呼び出しなど時間のかかる処理をバックグラウンドで非同期に実行できます。
APIサーバー Redis Worker
↓ job追加 → [Queue] → job取得・処理
↓ 完了/失敗
retry・遅延実行・優先度・並列数制御などが組み込まれており、TypeScriptファーストで設計されています。
インストール
npm install bullmq
RedisはDockerで起動します。
docker run -d --name redis -p 6379:6379 redis:latest
基本:Queue と Worker
BullMQの2つの主役は Queue(ジョブを追加する側)と Worker(ジョブを処理する側)です。
Queue:ジョブを追加する
// queue.ts import { Queue } from 'bullmq' const emailQueue = new Queue('email', { connection: { host: 'localhost', port: 6379 }, }) // ジョブを追加 await emailQueue.add('sendWelcome', { to: 'user@example.com', name: '田中太郎', }) console.log('ジョブを追加しました') await emailQueue.close()
Worker:ジョブを処理する
// worker.ts import { Worker, Job } from 'bullmq' type EmailJobData = { to: string name: string } const worker = new Worker<EmailJobData>( 'email', async (job: Job<EmailJobData>) => { console.log(`メール送信: ${job.data.to} へ「${job.name}」`) // 実際のメール送信処理 await sendEmail(job.data.to, job.data.name) }, { connection: { host: 'localhost', port: 6379 } } ) worker.on('completed', (job) => { console.log(`完了: ${job.id}`) }) worker.on('failed', (job, err) => { console.error(`失敗: ${job?.id}`, err.message) }) async function sendEmail(to: string, name: string) { console.log(`送信先: ${to}, 宛名: ${name}`) }
retry と backoff
失敗時の自動リトライを設定できます。
await emailQueue.add( 'sendWelcome', { to: 'user@example.com', name: '田中太郎' }, { attempts: 3, // 最大3回試みる backoff: { type: 'exponential', // 指数バックオフ delay: 1000, // 初回: 1秒, 2回目: 2秒, 3回目: 4秒 }, } )
| backoff type | 内容 |
|---|---|
fixed |
常に同じ待機時間 |
exponential |
待機時間を指数的に増やす(推奨) |
Workerの処理内でエラーをthrowするとretryが発動します。
const worker = new Worker('email', async (job) => { const result = await callExternalAPI(job.data) if (!result.ok) { throw new Error(`API error: ${result.status}`) // ← retryが発動 } })
遅延実行
// 10分後に実行 await emailQueue.add( 'sendReminder', { to: 'user@example.com' }, { delay: 10 * 60 * 1000 } ) // 翌朝9時に実行(Unixミリ秒で指定) const tomorrow9am = new Date() tomorrow9am.setDate(tomorrow9am.getDate() + 1) tomorrow9am.setHours(9, 0, 0, 0) await emailQueue.add( 'sendMorningDigest', { userId: 42 }, { delay: tomorrow9am.getTime() - Date.now() } )
優先度キュー
priority を指定すると小さい値ほど先に処理されます。
// 優先度: 1(最高) await emailQueue.add('sendAlert', { message: '緊急通知' }, { priority: 1 }) // 優先度: 10(低) await emailQueue.add('sendNewsletter', { message: '週刊' }, { priority: 10 })
進捗報告
長時間ジョブの進捗をリアルタイムで追跡できます。
// Worker側:進捗を更新 const worker = new Worker('imageProcess', async (job) => { await job.updateProgress(0) await downloadImage(job.data.url) await job.updateProgress(30) await resizeImage() await job.updateProgress(70) await uploadToS3() await job.updateProgress(100) }) // Queue側:進捗を受け取る const queueEvents = new QueueEvents('imageProcess', { connection: { host: 'localhost', port: 6379 }, }) queueEvents.on('progress', ({ jobId, data }) => { console.log(`Job ${jobId}: ${data}%`) })
ジョブのライフサイクルイベント
import { QueueEvents } from 'bullmq' const queueEvents = new QueueEvents('email', { connection: { host: 'localhost', port: 6379 }, }) queueEvents.on('waiting', ({ jobId }) => console.log(`待機中: ${jobId}`)) queueEvents.on('active', ({ jobId }) => console.log(`処理中: ${jobId}`)) queueEvents.on('completed', ({ jobId, returnvalue }) => console.log(`完了: ${jobId}`)) queueEvents.on('failed', ({ jobId, failedReason }) => console.error(`失敗: ${jobId} - ${failedReason}`))
管理UI(Bull Board)
npm install @bull-board/express @bull-board/api
import express from 'express' import { createBullBoard } from '@bull-board/api' import { BullMQAdapter } from '@bull-board/api/bullMQAdapter' import { ExpressAdapter } from '@bull-board/express' const serverAdapter = new ExpressAdapter() serverAdapter.setBasePath('/admin/queues') createBullBoard({ queues: [new BullMQAdapter(emailQueue)], serverAdapter, }) const app = express() app.use('/admin/queues', serverAdapter.getRouter()) app.listen(3000) // http://localhost:3000/admin/queues でキューの状態を確認できる
まとめ
| 機能 | 設定・メソッド |
|---|---|
| ジョブ追加 | queue.add(name, data) |
| retry | { attempts: 3, backoff: { type: 'exponential', delay: 1000 } } |
| 遅延実行 | { delay: ms } |
| 優先度 | { priority: 1 } (小さいほど優先) |
| 進捗報告 | job.updateProgress(percent) |
| 管理UI | Bull Board |
- PublisherとWorkerは別プロセスで動かすのが基本
Workerの処理内でthrowするとretryが発動するQueueEventsでキュー全体のイベントを監視できる
Redisを使ったPub/Subについては「Redis Pub/Sub入門:Node.jsとTypeScriptで実装する」も参照してください。