Redis
 Computer >> コンピューター >  >> プログラミング >> Redis

Node.jsとRedisで学ぶメッセージキュー入門:Webアプリのパフォーマンスを劇的に改善する方法

大規模なWebアプリケーションを開発する際、処理速度は最優先事項のひとつです。現代のユーザーは長いレスポンス待ちを許容しませんし、そもそも待たされるべきでもありません。しかし、どうしても時間がかかる処理には、高速化も省略もできないものが存在します。

そんな課題を解決してくれるのがメッセージキューです。通常の「リクエスト→レスポンス」の流れに別の処理経路を追加することで、ユーザーには即座にレスポンスを返しつつ、時間のかかる重い処理は裏側で実行できます。結果として、ユーザーにもシステムにも嬉しい仕組みが実現します。

本記事では、メッセージキューの基本概念を解説し、実際にシンプルなアプリケーションを構築しながらその使い方を学びます。Node.jsの基礎知識があること、およびRedisがローカルまたはクラウド環境にインストール済みであることを前提としています。

キュー(Queue)とは何か?

キューとは、データを特定の順序で格納できるデータ構造です。キューは「FIFO(First-In-First-Out:先入れ先出し)」という原則に従って動作します。

コンピュータサイエンスにおけるキューの概念は、日常生活での行列とまったく同じです。列の後ろに並び、自分の番が来るのを待ち、対応を受け終えたら前の方から抜けていく——それがキューの基本的なイメージです。

コンピュータサイエンスの世界では、APIリクエストのような処理を実行中に、「メール送信」などのタスクを現在のフローから切り離したい場合があります。そのようなとき、タスクをキューに積んで(プッシュして)、元の処理を続行すればよいのです。

以下の図は、キューのライフサイクルを示しています。

キューのライフサイクル | https://optimalbits.github.io/bull/

ジョブ(Job)とは何か?

ジョブとは、キュー上で扱われるデータの単位であり、一般的にはJSON形式のようなオブジェクトです。

空港で並んでいる乗客一人ひとりをジョブだと考えてみてください。それぞれの人が、固有のデータや必要な書類(パスポートや医療証明書など)が入ったブリーフケースを持っています。これらは自分の順番が来たときに対応してもらうための情報です。

新しく並ぶ人は列の後ろ(最後尾)に加わり、対応は前(先頭)から行われます。ジョブもまったく同じように処理されます。各ジョブは自身の処理に使われるデータを保持しており、新しいジョブは後ろから追加され、処理対象は前から取り出されます。

ジョブプロデューサー(Job Producer)とは?

ジョブプロデューサーとは、キューにジョブを追加するコードのことです。現実世界でいえば、空港で目的別にどの列に並ぶべきか案内する警備員のような存在です。

ジョブプロデューサーは、ジョブコンシューマー(消費者)から独立して存在できます。つまりマイクロサービス構成では、あるサービスはキューへのジョブ追加だけを担当し、その後の処理方法については関与しない、という分離が可能です。

ワーカー(Worker / ジョブコンシューマー)とは?

ワーカー(ジョブコンシューマー)とは、ジョブを実行できるプロセスまたは関数のことです。銀行の窓口係を想像してください。最初の客がやってきて列に並び、窓口係が呼び出すと列は空になります。

窓口係は取引処理に必要な詳細情報を客に求めます。その対応中にさらに4人の客が並んだとしても、彼らは最初の客の対応が終わるまで列で待機し、その後に次々と呼び出されていきます。キューのワーカーも同様に、キューの先頭にあるジョブを取り出しては処理していきます。

失敗したジョブ(Failed Jobs)とは?

実際の運用では、処理中に失敗するジョブが発生することがあります。

ジョブが失敗する主な理由は以下のとおりです。

  • 無効または欠落している入力データ: 処理に必要なデータが不足しているとジョブは失敗します。たとえば、宛先のメールアドレスがないメール送信ジョブは必ず失敗します。
  • タイムアウト: 通常より長くかかっているジョブは、キューメカニズムによって失敗扱いにされることがあります。依存先の不具合などが原因ですが、1つのジョブが永久に動き続けるのは望ましくありません。
  • ネットワークやインフラの問題: ほぼ制御不能な問題ですが、実際に起こりえます。たとえばデータベース接続エラーはジョブを失敗させます。
  • 依存関係の問題: ジョブが外部リソースに依存している場合、そのリソースが利用不可だったり処理に失敗したりすると、ジョブ自体も失敗します。

ジョブが失敗した場合のために、キューの仕組みにはリトライ(再試行)を設定できます。即時に再試行するか、計算された間隔をおいて再試行するかを選択でき、最大試行回数を設定しておくことが推奨されます。設定しないと、永遠に成功しないジョブを無限に実行し続けることになりかねません。

キューは、マイクロサービス間の堅牢な通信チャネルを構築するためにも役立ちます。複数のサービスが同じキューを共有でき、異なるサービスが異なる役割を担うことが可能です。あるサービスがタスクを完了したら、次の処理を待つワーカーがいる別のサービスへジョブをプッシュすれば、そのサービスがデータを受け取って必要な処理を行います。

また、キューはプロセスから重いタスクを切り離す用途にも有効です。本記事で見ていくように、メール送信のような時間のかかるタスクをキューに委ねることで、レスポンスタイムの低下を防げます。

さらに、キューは単一障害点(SPOF)の回避にも貢献します。失敗の可能性があり、再試行が必要な処理は、少し時間をおいて再実行できるキューで扱うのが最適です。

キューを使ったシンプルなアプリの作り方

ここからは、Node.jsとRedisを使った簡単なプロジェクトを構築します。キューシステム構築に伴う複雑さを大幅に軽減してくれるBullライブラリを採用します。このプロジェクトは、メール送信用の単一エンドポイントを持つアプリです。

新しいNode.jsプロジェクトを作成する

mkdir nodejs-queue-project
cd nodejs-queue-project
npm init -y

上記のコマンドにより、nodejs-queue-projectという名前のフォルダと、その中にpackage.jsonファイルが作成されます。package.jsonは次のようになっているはずです。

{
 "name": "nodejs-queue-project",
 "version": "1.0.0",
 "description": "",
 "main": "index.js",
 "scripts": {
 "test": "echo \"Error: no test specified\" && exit 1"
 },
 "keywords": [],
 "author": "",
 "license": "ISC"
}

必要な依存パッケージをインストールする

npm i express @types/express @types/node body-parser ts-node ts-lint typescript nodemon nodemailer @types/nodemailer

上記のコマンドで、プロジェクトに必要な各種パッケージと依存関係がインストールされます。

インストール後、package.jsonscriptsセクションにdevコマンドを追加しましょう。package.json全体は次のようになります。

{
 "name": "nodejs-queue-project",
 "version": "1.0.0",
 "description": "",
 "main": "index.js",
 "scripts": {
 "dev": "nodemon src/app.ts"
 },
 "keywords": [],
 "author": "",
 "license": "ISC",
 "dependencies": {
 "@types/express": "^4.17.17",
 "@types/node": "^20.3.3",
 "@types/nodemailer": "^6.4.8",
 "body-parser": "^1.20.2",
 "express": "^4.18.2",
 "nodemailer": "^6.9.3",
 "nodemon": "^2.0.22",
 "ts-lint": "^4.5.1",
 "ts-node": "^10.9.1",
 "typescript": "^5.1.6"
 }
}

このファイルにはインストール済みのすべての依存関係が記録されています。devスクリプトを使うと、npm run devコマンドでアプリが起動します。

エンドポイントの構築方法

まず、srcという名前の新しいフォルダを作成します。このフォルダにすべてのコードファイルを格納します。最初のファイルはアプリケーションのルートファイル、つまりpackage.jsonで定義したapp.tsです。

app.tsでは、必要なパッケージをインポートし、メール送信用の単一エンドポイントを持つシンプルなサーバーを以下のように作成します。

import express from "express";
import bodyParser from "body-parser";
import nodemailer from "nodemailer";
const app = express();
app.use(bodyParser.json());
app.post("/send-email", async (req, res) => {
 const { from, to, subject, text } = req.body;
 // これはチュートリアルなのでテストアカウントを使用
 const testAccount = await nodemailer.createTestAccount();
 const transporter = nodemailer.createTransport({
 host: "smtp.ethereal.email",
 port: 587,
 secure: false,
 auth: {
 user: testAccount.user,
 pass: testAccount.pass,
 },
 tls: {
 rejectUnauthorized: false,
 },
 });
 console.log("Sending mail to %s", to);
 let info = await transporter.sendMail({
 from,
 to,
 subject,
 text,
 html: `<strong>${text}</strong>`,
 });
 console.log("Message sent: %s", info.messageId);
 console.log("Preview URL: %s", nodemailer.getTestMessageUrl(info));
 res.json({
 message: "Email Sent",
 });
});
app.listen(4300, () => {
 console.log("Server started at https://localhost:4300");
});

ターミナルでnpm run devを実行するとサーバーが起動し、Server started at https://localhost:4300というメッセージが表示されるはずです。

次に、Postmanなどのツールでエンドポイントをテストしてみましょう。

Postmanによるエンドポイントのテスト

スクリーンショットのとおり、このリクエストには約4秒かかりました。エンドポイントとしては非常に遅い値です。ターミナルを確認すると、送信されたメールをプレビューできるURLも表示されているはずです。

リンクを開くと、送信されたメールの内容を確認できます。

メールの内容

キューの作成方法

処理をさらに高速化するために、メール送信を後回しにしてキューに登録し、ユーザーにはすぐにレスポンスを返すようにします。

そのためには、キュー作成に使用するbullライブラリとその型定義パッケージ@types/bullをインストールします。

npm i bull @types/bull

bullを使った新しいキューの作成は非常に簡単で、キュー名を指定してBullオブジェクトをインスタンス化するだけです。

// ファイルの先頭に記述します
import Bull from 'bull';
const emailQueue = new Bull("email");

キュー名だけで作成した場合、デフォルトのRedis接続URLであるlocalhost:6379が使用されます。別のURLを使いたい場合は、第2引数としてオプションオブジェクトを渡します。

const emailQueue = new Bull("email", {
 redis: "localhost:6379",
});

ここまで準備できたら、ジョブプロデューサーとして機能するシンプルな関数を作成し、リクエストが届くたびにジョブをキューへ追加できるようにしましょう。

type EmailType = {
 from: string;
 to: string;
 subject: string;
 text: string;
};
const sendNewEmail = async (email: EmailType) => {
 emailQueue.add({ ...email });
};

新しく作成した関数sendNewEmailは、EmailType型のオブジェクト(送信元アドレスfrom、宛先アドレスto、件名subject、本文text)を受け取り、新しいジョブをキューにプッシュします。

次に、リクエスト内で直接メールを送信する代わりに、この関数を使うようエンドポイントを修正します。

app.post("/send-email", async (req, res) => {
 const { from, to, subject, text } = req.body;
 await sendNewEmail({ from, to, subject, text });
 console.log("Added to queue");
 res.json({
 message: "Email Sent",
 });
});

これでコードはシンプルになり、処理も大幅に速くなりました。リクエストにかかる時間は約40ミリ秒——以前と比べて約100倍の高速化です。

Postmanによるエンドポイントのテスト

この時点で、メールはキューに追加されています。処理されるまでキューの中で待機し続けます。ジョブの処理は同じアプリケーションでも、別のサービス(マイクロサービス構成の場合)でも実行できます。

ジョブの処理方法

メールがキューから取り出されないままでは、この仕組みは不完全で意味がありません。そこで、ジョブを処理してキューを空にするジョブコンシューマーを作成します。

具体的には、Jobオブジェクトを受け取ってメールを送信する関数のロジックを作成します。

const processEmailQueue = async (job: Job) => {
 // これはチュートリアルなのでテストアカウントを使用
 const testAccount = await nodemailer.createTestAccount();
 const transporter = nodemailer.createTransport({
 host: "smtp.ethereal.email",
 port: 587,
 secure: false,
 auth: {
 user: testAccount.user,
 pass: testAccount.pass,
 },
 tls: {
 rejectUnauthorized: false,
 },
 });
 const { from, to, subject, text } = job.data;
 console.log("Sending mail to %s", to);
 let info = await transporter.sendMail({
 from,
 to,
 subject,
 text,
 html: `<strong>${text}</strong>`,
 });
 console.log("Message sent: %s", info.messageId);
 console.log("Preview URL: %s", nodemailer.getTestMessageUrl(info));
 return nodemailer.getTestMessageUrl(info);
};

この関数はJobオブジェクトを受け取ります。Jobオブジェクトには、ジョブのステータスやデータを確認できる便利なプロパティが含まれており、ここではdataプロパティを使用しています。

ただし、現状ではこれは単なる関数です。どのキューを担当すべきかを知らないため、自動的にはジョブを拾いません。

キューに接続する前に、いくつかリクエストを送信してジョブをキューに追加してみましょう。現在待機中のメールジョブは、redis-cliで次のコマンドを実行すると確認できます。

LRANGE bull:email:wait 0 -1

このコマンドはメールの待機リストをチェックし、待機中ジョブのidを返します。

Redis CLI

ワーカーの実際の挙動を確認するため、いくつかジョブを作成しておきました。

それでは、次の1行のコードを追加して、ワーカーをキューに接続します。

emailQueue.process(processEmailQueue);

追加後のapp.tsファイル全体は次のようになります。

import express from "express";
import bodyParser from "body-parser";
import nodemailer from "nodemailer";
import Bull, { Job } from "bull";
const app = express();
app.use(bodyParser.json());
const emailQueue = new Bull("email", {
 redis: "localhost:6379",
});
type EmailType = {
 from: string;
 to: string;
 subject: string;
 text: string;
};
const sendNewEmail = async (email: EmailType) => {
 emailQueue.add({ ...email });
};
const processEmailQueue = async (job: Job) => {
 // これはチュートリアルなのでテストアカウントを使用
 const testAccount = await nodemailer.createTestAccount();
 const transporter = nodemailer.createTransport({
 host: "smtp.ethereal.email",
 port: 587,
 secure: false,
 auth: {
 user: testAccount.user,
 pass: testAccount.pass,
 },
 tls: {
 rejectUnauthorized: false,
 },
 });
 const { from, to, subject, text } = job.data;
 console.log("Sending mail to %s", to);
 let info = await transporter.sendMail({
 from,
 to,
 subject,
 text,
 html: `<strong>${text}</strong>`,
 });
 console.log("Message sent: %s", info.messageId);
 console.log("Preview URL: %s", nodemailer.getTestMessageUrl(info));
};
emailQueue.process(processEmailQueue);
app.post("/send-email", async (req, res) => {
 const { from, to, subject, text } = req.body;
 await sendNewEmail({ from, to, subject, text });
 console.log("Added to queue");
 res.json({
 message: "Email Sent",
 });
});
app.listen(4300, () => {
 console.log("Server started at https://localhost:4300");
});

保存すると、サーバーが再起動し、すぐにメールの送信が始まるのがわかるでしょう。ワーカーがキューを検知して、即座に処理を開始するためです。

サーバーが待機中のメールを送信している様子

これでプロデューサーとワーカーの両方が稼働状態になりました。新しいAPIリクエストはすべてキューにプッシュされ、保留中のジョブがない限り、ワーカーが即座に処理します。

まとめ

本記事を通じて、メッセージキューとは何か、ジョブの追加方法と実行プロセスの作成方法、そしてそれらを活用してより良いWebアプリケーションを構築する方法を理解できたなら幸いです。記事で使用したコードファイルはGitHubで公開しています。

ご質問やフィードバックがあれば、ぜひお気軽にお寄せください。


  1. データベースレス(DBLess)アーキテクチャとは?ファーストプリンシプル思考で読み解くデータベースの未来

    「なぜRedisのようなデータベース企業が、Databaseless(DBLess:データベースレス)アーキテクチャについて語るのだろう?」——そう疑問に感じるのは自然なことです。まずは細部に入る前に、このまったく新しいアーキテクチャの背景にある「新しい考え方」から見ていきましょう。 そのために、少しだけ回り道をして「ファーストプリンシプル(第一原理)思考」についてお話しします。この思考法は、伝統や慣習を鵜呑みにせず、自らの頭で考え、あらゆるものに疑問を投げかけることを促すものです。 ファーストプリンシプルの核心はこうです。重力の法則のような自然法則でない限り、あらゆるシステムや概念は人間

  2. Redis ZREVRANGEBYLEXコマンドの使い方|辞書順降順で範囲指定してソート済みセットの要素を取得する方法

    このチュートリアルでは、RedisのZREVRANGEBYLEXコマンドを使って、ソート済みセット(Sorted Set)から特定の範囲内の値を持つすべての要素を、辞書順の降順で取得する方法を解説します。 ZREVRANGEBYLEXコマンドとは ZREVRANGEBYLEXコマンドは、指定したキーに保存されているソート済みセットの要素のうち、max引数とmin引数で指定された範囲(要素の文字列表現)に該当するすべての要素を返します。このコマンドを利用する際は、ソート済みセットのすべての要素を同じスコアで挿入し、辞書順(レキシコグラフィカル順)での並び替えを強制するのが前提となります。返される