BullMQとは:Node.js(TypeScript)のジョブキューライブラリ入門

スポンサーリンク

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で実装する」も参照してください。