📢

NestJS + Bullで学ぶメッセージキュー入門

に公開

はじめに

なぜこの記事を書いたか

先日、担当しているサービスで「100人以上の送信対象者に一括でメールを送信する」機能を実装しました。
メール送信処理で、Bull(メッセージキュー)を使いましたが、正直に言うと今までBullを何となく便利なものとして使っていて、Bullの仕組みや役割を深く理解していませんでした。

この記事では、自分自身の学び直しも兼ねて、メッセージキューの基本から、NestJS + Bullによる100人以上の送信対象者へのメール送信の実装までをまとめようと思います。

1. Bullを使わない同期処理の限界

まず、Bullを使わずに同期処理で100件のメール送信を行うとどうなるか考えてみます。

// 同期処理の例
@Mutation(() => Boolean)
async sendBulkMessage(
@Args("input") input: SendBulkMessageInput,
) {
  const users = await this.userRepository.findByShopId(input.shopId);

  // 100件のメール送信を直列で実行
  for (const user of users) {
    await this.sendgridService.send({
      to: user.email,
      subject: input.title,
      text: input.body,
    });
  }

  return true;
}

上記のコードの問題点

  • タイムアウト
    • 1件100msでも100件で10秒。HTTPタイムアウトの危険がある。
  • 全体が失敗する
    • 1件のエラーで処理全体が止まるリスクがある。
  • UX悪化
    • ユーザーは送信完了まで画面を待ち続けることになる。

これらの問題を解決するのがメッセージキューです。

2. メッセージキュー・Bullとは

メッセージキューとは

メッセージキューは、処理の依頼(メッセージ)を一時的に蓄積する仕組みです。
蓄積されたメッセージは、ワーカーと呼ばれる処理担当が順番に取り出して実行します。

Bullとは

Bullは、Node.js向けのメッセージキューライブラリです。
https://optimalbits.github.io/bull/
(※現在Bullはバグ修正のみ対応されているため、BullMQの使用・移行が推奨されています。)

Bullの主な特徴

  • Redisをジョブの保存・管理に使用
  • リトライ、遅延実行、優先度付けなどをサポート
  • TypeScript対応
  • 2014年から開発されている成熟したライブラリ

Bullが解決する課題

課題 Bullによる解決
非同期化 APIは即座にレスポンスを返し、処理はバックグラウンドで実行される
リトライ 失敗したジョブを自動で再試行する(回数・間隔を設定可能)
負荷分散 複数のワーカーでジョブを並列処理する
永続化 Redisにジョブが保存され、サーバー再起動後も復旧可能
監視 ジョブの状態(待機・処理中・完了・失敗)を追跡可能

BullがRedisを使う理由
Bullはジョブの保存先としてRedisを必須で使用します。
Redisの特性がキューに最適な理由は以下の通りです。

Redisの特性 キューでのメリット
インメモリ ジョブの追加・取得が高速(ミリ秒単位)
永続化オプション サーバー再起動後もジョブを復旧可能
アトミック操作 ロックにより同時に同じジョブを処理することを防ぐ(ただしat least onceのため、stalled後の再実行などで重複処理の可能性があり、冪等性の担保が必要)
Pub/Sub ジョブの状態変更をリアルタイムに通知
TTL(有効期限) Bullのオプション(removeOnComplete、removeOnFail等)と組み合わせて、完了/失敗ジョブの自動削除に活用可能

もしRedisを使わなかった場合、以下の問題が発生します。

  • メモリ内配列でキュー管理 → サーバー再起動でジョブが消える
  • ファイルベース → 複数サーバーで共有できない
  • RDBMSでキュー実装 → ポーリングが必要で非効率

@nestjs/bull と bull の違い

担当サービスでは @nestjs/bullbull の2つのパッケージを使用しています。それぞれの役割は以下の通りです。

パッケージ 役割
bull メッセージキューのコア機能を提供するライブラリ本体
@nestjs/bull bullをNestJSで使いやすくするラッパー(デコレーターやDIを提供する)

https://docs.nestjs.com/techniques/queues

// bull単体で使う場合
import Queue from "bull";
const myQueue = new Queue("mail", "redis://127.0.0.1:6379");
myQueue.add({ email: "test@example.com" });

// @nestjs/bullを使う場合(NestJSのDIに統合される)
@Injectable()
export class MailService {
  constructor(@InjectQueue("mail") private mailQueue: Queue) {}

  async addMail(data: MailInput) {
    await this.mailQueue.add("sendMail", data);
  }
}

@nestjs/bullを使うことで、以下のことが可能になります。

  • @InjectQueue()でキューをDI(依存性注入)できる
  • @Processor() / @Process()デコレーターでワーカーを定義できる
  • NestJSのモジュールシステムと統合できる

3. いつキューを使うべきか?

キューが有効なケース

  1. 外部API呼び出しを伴う処理
  • メール送信(SendGrid, SES)
  • LINE / Slack通知
  1. 処理時間が不安定・長い
  • ネットワーク遅延の影響を受ける
  • レート制限がある外部サービス
  1. 失敗時にリトライが必要
  • 一時的なネットワークエラー
  • 外部サービスの一時障害
  1. 大量の処理を並列化したい
  • 一括メール送信
  • バッチ処理
  1. ユーザーに即座にレスポンスを返したい
  • フォーム送信後すぐに完了画面を表示し、メール送信自体は裏側で非同期実行

キューが過剰なケース

  1. 数ms〜数十msで完了する単発のDB更新
  2. 決済結果の即時表示など同期的なフィードバックが必須な処理
  3. 処理件数が少なく、失敗しても手動リカバリで十分な処理
  4. Redis等のインフラ追加コストに見合わない規模のプロジェクト

判断フローチャート

4. [実践] NestJSでのモジュール設計

一括で複数のメール送信する機能の実装では、 Producer(キュー登録)Consumer(ワーカー) を分離した、プロジェクトの既存のディレクトリ構成をベースに解説します。(実際のプロジェクトでは少し異なる構成になっていますが、本記事では以下の構成とします。)

ディレクトリ構成

src/        
├── queue/          # Producer(キューへの登録)    
│   └── mail/
│       ├── mail-queue.module.ts
│       └── mail-queue.service.ts
└── processor/      # Consumer(ジョブの処理)
  └── mail/
      ├── mail-processor.module.ts
      └── mail.processor.ts

キューモジュール(Producer側)

ジョブをキューに登録する側のモジュールは以下のようになっています。

モジュール定義

// src/queue/mail/mail-queue.module.ts
import { Module } from "@nestjs/common";
import { BullModule } from "@nestjs/bull";
import { MailQueueService } from "./mail-queue.service";       
@Module({
  imports: [
    BullModule.registerQueue({
      name: "mail",  
      redis: {
        host: process.env.REDIS_HOST,
        port: parseInt(process.env.REDIS_PORT),    
      },
    }),
  ],
  providers: [MailQueueService],
  exports: [MailQueueService],  // 他モジュールから使えるようにexport
})
export class MailQueueModule {}

ジョブデータの型定義とサービス

// src/queue/mail/mail-queue.service.ts
import { Injectable } from "@nestjs/common";
import { InjectQueue } from "@nestjs/bull";
import { Queue } from "bull";
export interface SendMailJobData {
  to: string;
  subject: string;
  body: string;
}
@Injectable()
export class MailQueueService {
  constructor(
    @InjectQueue("mail") private readonly mailQueue: Queue,
  ) {}
  async addSendMailJob(data: SendMailJobData): Promise<void> {
    await this.mailQueue.add("sendMail", data, {
      attempts: 5,              // 最大5回試行(初回1回 + リトライ4回)
      backoff: {
        type: "exponential",    // 指数バックオフ
        delay: 5000,            // 初回5秒後にリトライ
      },
    });
  }
}

ポイント

  • SendMailJobDataインターフェースでジョブデータの型を定義し、Producer/Consumer間で共有
  • @InjectQueue("mail")でキューをDI
  • add()の第1引数がジョブ名、第2引数がデータ、第3引数がオプション

リトライ戦略(指数バックオフ)
上記のbackoff設定により、以下のタイミングでリトライが実行されます。

試行 タイミング 計算式
初回 即時実行 -
リトライ1回目 5秒後 delay: 5000
リトライ2回目 10秒後 5000 × 2^1
リトライ3回目 20秒後 5000 × 2^2
リトライ4回目 40秒後 5000 × 2^3

外部サービスが一時的に障害状態のとき、すぐにリトライすると負荷をかけてしまいますが、指数バックオフで時間を空けることで、サービス復旧を待つ余裕が生まれます。

プロセッサーモジュール(Consumer側)

次に、キューからジョブを取り出して処理する側のモジュールは以下のようになっています。
モジュール定義

// src/processor/mail/mail-processor.module.ts
import { Module } from "@nestjs/common";
import { BullModule } from "@nestjs/bull";
import { MailProcessor } from "./mail.processor";
import { MailModule } from "@/mail/mail.module";
@Module({
  imports: [
    BullModule.registerQueue({
      name: "mail",
      redis: {
        host: process.env.REDIS_HOST,
        port: parseInt(process.env.REDIS_PORT),
      },
    }),
    MailModule,  // SendGridなどのメール送信サービス
  ],
  providers: [MailProcessor],       
})
export class MailProcessorModule {}

プロセッサー(ワーカー)

// src/processor/mail/mail.processor.ts
import { Process, Processor } from "@nestjs/bull";
import { Injectable } from "@nestjs/common";
import { Job } from "bull";
import { MailService } from "@/mail/mail.service";
import { SendMailJobData } from "@/queue/mail/mail-queue.service";
@Processor("mail")
@Injectable()
export class MailProcessor {
  constructor(private readonly mailService: MailService) {}
  @Process("sendMail")
  async handleSendMail(job: Job<SendMailJobData>): Promise<void> {
    const { to, subject, body } = job.data;
    await this.mailService.send({     
      to,  
      subject,
      text: body,
    });
  }
}

ポイント

  • @Processor("mail")でどのキューを処理するか指定する
  • @Process("sendMail")で処理するジョブ名を指定する
  • job.dataでProducer側が登録したデータを取得する

アプリケーションモジュールへの登録

// src/app.module.ts
@Module({
  imports: [
    // ... 他のモジュール
    MailQueueModule,       // Producer
    MailProcessorModule,   // Consumer
  ],
})
export class AppModule {}

5. 一括メール送信機能を作る

セクション4のモジュール構成を使って、一括メール送信機能を実装します。

// bulk-message.usecase.ts
import { Injectable } from "@nestjs/common";
import { MailQueueService } from "@/queue/mail/mail-queue.service";
@Injectable()
export class BulkMessageUsecase {
  constructor(
    private readonly mailQueueService: MailQueueService,
    private readonly userRepository: UserRepository,
  ) {}
  async sendBulkMessage(input: SendBulkMessageInput) {
    const users = await this.userRepository.findByShopId(input.shopId);
    // 各ユーザーへの送信をキューに登録
    for (const user of users) {
      await this.mailQueueService.addSendMailJob({
        to: user.email,
        subject: input.title,
        body: input.body,
      });  
    }
    // キューへの登録が完了したらレスポンス
    // 実際のメール送信はバックグラウンドで MailProcessor が実行する
    return true;
  }
}

実際のメール送信はバックグラウンドでMailProcessorが実行するため、SendGridへのAPI呼び出しを待つ必要はありません。しかし、この実装にはいくつかの課題があります。

課題1: 送信対象者数に比例するレスポンス時間
usecase内でforループを回しているため、送信対象者数に比例してRedis操作が増えます。

送信対象者数 APIリクエスト内のRedis操作 レスポンス時間
10人 10回 数十ms
100人 100回 数百ms
1000人 1000回 数秒

Redis操作は1回あたり数ミリ秒程度ですが、送信対象者が増えると体感できる遅延になります。

課題2: 送信結果がわからない
100件中何件成功したか追跡できません。

課題3: 後処理ができない
APIはキュー登録完了時点でレスポンスを返すため、「全件送信完了後に結果通知を送る」といった後処理ができません。

次のセクションで、これらの課題を解決する2段階キュー構造を作成します。

6. 2段階キュー構造と送信結果通知

2段階キュー構造

セクション5の課題を解決するため、2段階キュー構造を採用しました。

[APIリクエスト]
  └─ bulk-message キューにジョブ追加(1回)
  └─ レスポンス返却

[bulk-message processor(非同期)]
  └─ for (送信対象者)
        └─ mail キューにジョブ追加
  └─ 結果通知メールを mail キューに追加

2段階構造による解決

セクション5の課題 2段階構造での解決方法
課題1: レスポンス時間 親ジョブ1つだけ登録(Redis操作1回)で即座にレスポンス
課題2: 送信結果 親ジョブのProcessor内で成功/失敗を集計
課題3: 後処理 親ジョブのProcessor内で全件処理後に結果通知を送信

forループをバックグラウンドの親ジョブに移動することで、送信対象者数に関係なく即座にレスポンスを返せます。

送信対象者数 1段階(Redis操作N回) 2段階(Redis操作1回)
10人 数十ms 数ms
100人 数百ms 数ms
1000人 数秒 数ms

2段階構造の設計
[親ジョブ: bulk-message キュー]

  • 一括送信の「単位」を管理
  • 送信対象リストを保持
  • 全メール送信完了後に結果通知を送信

[子ジョブ: mail キュー]

  • 1通ずつのメール送信を管理
  • 個別のリトライ
  • SendGridとの通信

実装例

親ジョブのデータ構造

// src/queue/bulk-message/bulk-message-queue.service.ts
export interface SendBulkMessageJobData {
  bulkMessageId: number;
  senderName: string;
  notificationEmail: string;   // 結果通知の送信先
  title: string;
  body: string;
  emailTargets: Array<{        // メール送信対象リスト
    email: string;
    userName: string;
  }>;
}

@Injectable()
export class BulkMessageQueueService {
  constructor(
    @InjectQueue("bulk-message") private readonly bulkMessageQueue: Queue,
  ) {}
  async addSendEmailJob(data: SendBulkMessageJobData): Promise<void> {
    await this.bulkMessageQueue.add("sendEmails", data, {
      attempts: 3,
      backoff: {
        type: "exponential",
        delay: 5000,
      },
    });
  }
}

親ジョブのProcessor

// src/processor/bulk-message/bulk-message.processor.ts
@Processor("bulk-message")
export class BulkMessageProcessor {
  constructor(
    private readonly mailService: MailService,
  ) {}
  @Process("sendEmails")
  async handleSendEmails(job: Job<SendBulkMessageJobData>): Promise<void> {
    const { senderName, notificationEmail, title, body, emailTargets } = job.data;
    let successCount = 0;
    const failedUsers: Array<{ name: string }> = [];
    // 複数の宛先にメール送信
    for (const target of emailTargets) {
      try {
        // 内部で mail キューにジョブを登録
        // N回呼び出しはバックグラウンドで実行されるため、操作者の待ち時間には影響しない
        await this.mailService.sendBulkMessageNotification(
          target.email,
          senderName,
          title,
          body,
        );
        successCount++;
      } catch (error) {
        failedUsers.push({ name: target.userName });
      }
     }
     // 全件処理完了後、結果通知メールを送信
     // これも内部で mail キューにジョブを登録
     await this.mailService.sendBulkMessageResultNotification(
       notificationEmail,   
       title,
       new Date(),
       successCount,
       failedUsers.length,
       emailTargets.map((t) => ({ name: t.userName })),
       failedUsers,
     );
   }
}

ポイント

  • mailService.sendBulkMessageNotification()は内部でmailキューにジョブを登録する
  • 親ジョブ内のforループは「mailキューへの登録」であり、実際のSendGrid通信ではない
  • mailキューへの登録成功を「送信成功」として集計
  • 注: 実際のSendGrid APIの結果ではなく、キュー登録の成功を集計している

Usecaseの実装

// bulk-message.usecase.ts
@Injectable()
export class BulkMessageUsecase {
  constructor(
    private readonly bulkMessageQueueService: BulkMessageQueueService,
    private readonly userRepository: UserRepository,
  ) {}
  async sendBulkMessage(input: SendBulkMessageInput) {
    const users = await this.userRepository.findByShopId(input.shopId);
    // 親ジョブを1つだけキューに登録(Redis操作1回)
    await this.bulkMessageQueueService.addSendEmailJob({
      bulkMessageId: input.id,
      senderName: input.senderName,
      notificationEmail: input.notificationEmail,
      title: input.title,
      body: input.body,
      emailTargets: users.map((u) => ({
        email: u.email,
        userName: u.name,
      })),
    });
    // 即座にレスポンス
    return true;
  }
}

いつ2段階構造を使うべきか?
2段階構造が有効なケース

  • 一括処理の完了後に後処理(通知、集計)が必要
  • 送信結果をユーザーに報告する必要がある
    1段階で十分なケース
  • 単純な一括送信で結果通知が不要
  • 各送信が完全に独立している

最後に

今回、一括メール送信機能の実装を通じて、NestJS + Bullによるメッセージキューの設計パターンを整理しました。
これまでメール送信を実装する際、Bullを「なんとなく便利なもの」として使っていましたが、2段階キュー構造を設計・実装したことで、以下の点を理解できました。

  • キューへの登録(Redis操作)と実際の処理(SendGrid通信)の違い
  • 親ジョブと子ジョブを分けることで得られる柔軟性
  • レスポンス時間と処理の信頼性のトレードオフ

「Bullを使っているけど仕組みがよくわからない」「一括処理の設計に悩んでいる」という方の参考になれば幸いです。

GMOメディアテックブログ

Discussion