RabbitMQ と Kafka
この部の 7 / 13 章 ・ 全体で 49 / 76 章 ・ 読了目安 45 分
- 「消えるか残るか」で方式を選べる
- prefetch と DLQ を設定できる
- パーティション数が並列度の上限だと理解している
前章では、非同期メッセージングの考え方を扱いました。 この章では、実際の製品を見ます。仕組みが分かると、設計判断ができるようになります。
実務で出会うのは、だいたいこの3種類です。
| 代表 | 一言でいうと | |
|---|---|---|
| メッセージブローカー | RabbitMQ, ActiveMQ | 仕事を配る。届けたら消える |
| ログ型ストリーム | Kafka, Pulsar | 記録を残す。読んでも消えない |
| マネージド | Cloud Pub/Sub, SQS/SNS | 運用を任せる。中身はどちらかに近い |
「消えるか、残るか」が最大の違いです。ここから他の性質がほぼ決まります。
もともとメッセージングは、企業システムの連携のために生まれました (メインフレームと業務システムをつなぐ、など)。 そこで求められたのは「確実に1回、正しい相手に届ける」ことでした。 この系譜が RabbitMQ などのブローカーです。
一方 2010 年代、Web サービスが扱うデータ量が爆発しました。 「大量のログやイベントを、複数のシステムがそれぞれのペースで読みたい」 という要求が出てきます。届けたら消す方式ではこれができません。
そこで LinkedIn が作ったのが Kafka で、発想が逆転しています。 メッセージを消さず、追記だけされるログとして保持し、 読む側が「どこまで読んだか」を自分で覚える。
用途が違うので、優劣ではありません。 どちらの系譜かを見分けてください。
RabbitMQ — 仕事を配る
まずはこちらから理解するのが分かりやすいです。
3つの登場人物
[送信者] --> [Exchange] --(Binding)--> [Queue] --> [受信者]
どこへ送るか判断 仕事の入れ物
Exchange 受け取ったメッセージを、どの Queue に入れるか決める
Queue メッセージが溜まる場所。ここから取り出される
Binding Exchange と Queue を結ぶルール(どんな条件で流すか)
送信者は Queue を知りません。 Exchange に投げるだけです。 これにより、後から受信側を追加しても送信側を変えずに済みます。
Exchange の種類
| 種類 | 振り分け方 | 使いどころ |
|---|---|---|
| direct | ルーティングキーが完全一致 | 「この種類の仕事はこのキューへ」 |
| fanout | 全部の Queue に複製 | 1つの出来事を複数のサービスに通知 |
| topic | パターン一致(order.*.jp) | 条件で細かく振り分ける |
| headers | ヘッダの内容で判定 | あまり使わない |
fanout の例:
[注文確定] → [Exchange(fanout)] → [在庫キュー]
→ [通知キュー]
→ [分析キュー]
これがドメイン駆動設計の実践のドメインイベントの実装形になります。
受け取りと ack
1. コンシューマがメッセージを受け取る
2. 処理する
3. ack(処理完了)を返す ← ここで初めてキューから消える
4. ack が返らなければ、別のコンシューマに再配信される
「受け取った時点」ではなく「ack した時点」で消える—— これが at-least-once の正体です(前章)。
□ ack は処理が完全に終わってから返す
□ 処理に失敗したら nack して、再試行させるか DLQ に送る
□ 自動 ack(受け取った瞬間に消える)は、失うことを許容できる場合だけ
既定では、コンシューマが大量のメッセージをまとめて抱え込むことがあります。
コンシューマA が100件を抱えて処理中
コンシューマB は暇なのに、何も回ってこない
prefetch = 1 1件ずつ。処理時間が長い仕事に向く
prefetch = 10〜 細かい仕事を高速に捌く場合
「増やしたのに速くならない」の原因になります。 処理時間が長いワーカーでは、まず prefetch を疑ってください。
DLQ の作り方
RabbitMQ では Dead Letter Exchange(DLX) を設定します。
1. Queue に「失敗したらこの Exchange へ」を設定する
2. nack された、期限切れになった、上限を超えたメッセージがそこへ流れる
3. DLQ に溜まったものを、人が調査して再投入する
前章のとおり、DLQ は「捨て場」ではなく「要調査の箱」です。 件数を監視してください。
「失敗したら同じキューに戻す」を素朴に実装すると、 永久に回り続けます(毒メッセージ。前章)。
□ リトライ回数をヘッダに持たせ、上限を超えたら DLQ へ
□ 遅延キュー(一定時間後に戻す)を使い、間隔を空ける
Kafka — 記録を残す
Kafka は、メッセージキューというより「追記専用のログ」です。
Topic: orders
Partition 0: [msg][msg][msg][msg][msg] ...
0 1 2 3 4 ← オフセット(何番目か)
Partition 1: [msg][msg][msg] ...
Partition 2: [msg][msg][msg][msg] ...
Topic メッセージの種類(RabbitMQ の Exchange に近い)
Partition Topic を分割したもの。**並列度の単位**
Offset 各パーティション内の通し番号
決定的に違う点
| RabbitMQ | Kafka | |
|---|---|---|
| 読んだら | 消える | 消えない(保持期間まで残る) |
| 進捗の管理 | ブローカーが持つ | コンシューマがオフセットを持つ |
| 読み直し | できない | できる(オフセットを戻す) |
| 順序 | 保証しにくい | パーティション内では保証される |
| 複数の購読者 | Exchange で複製する | 同じログを各自のペースで読む |
Kafka の最大の利点は、これです。
□ バグを直した後、過去のデータを最初から処理し直せる
□ 新しいサービスを追加して、過去分から取り込める
□ 障害で処理が飛んだ分を、オフセットを戻して復旧できる
RabbitMQ ではメッセージは消えているので、これができません。
逆に言えば、Kafka は保持期間の分だけディスクを使い、 運用の難易度も上がります。「読み直しが必要か」で判断してください。
コンシューマグループ
Topic のパーティション: P0, P1, P2
グループA(3台): 各台が1パーティションずつ担当 → 並列に処理
グループB(1台): 1台が全パーティションを担当 → 同じデータを独立に処理
□ 同じグループ内では、1つのパーティションを1台だけが読む
□ **コンシューマを増やしても、パーティション数を超えると余る**
□ 別のグループなら、同じデータを独立に読める(ファンアウト)
パーティション 3 に対して、コンシューマ 10台
→ 3台だけが働き、7台は何もしない
パーティション数は後から増やせますが、減らせません。 そして増やすと、キーとパーティションの対応が変わります(順序に影響)。
最初の設計で、想定する並列度より少し多めにしておくのが定石です。
順序とキー
同じキーのメッセージは、同じパーティションに入る
→ そのキーの中では順序が保証される
キー = 注文ID → 同じ注文の「作成 → 更新 → 削除」は順番通りに届く
キーなし → ラウンドロビンで分散。順序は保証されない
順序が必要な単位をキーにする——これが Kafka の設計の要点です。 ただし前章のとおり、順序に依存しない設計のほうが強いことは変わりません。
マネージドサービス
自前で運用するのは、それなりの負担です(クラスタ、ディスク、監視、バージョンアップ)。 クラウドのマネージドサービスを使うことが多くなっています。
| サービス | 近い方式 | 特徴 |
|---|---|---|
| Cloud Pub/Sub | 両者の中間 | トピック + サブスクリプション。運用が要らない |
| Amazon SQS | ブローカー型 | シンプルなキュー |
| Amazon SNS | fanout | 通知の配信 |
| Cloud Tasks | タスク実行 | 「この URL をこの時刻に叩く」 |
| Amazon MSK / Confluent | Kafka | Kafka そのものをマネージドで |
Google Cloud を使うなら、まずこれが選択肢になります。
Topic 送信先
Subscription 購読の単位。**トピックごとに複数作れる**
→ サブスクリプションを分ければファンアウトになる
Push / Pull サーバーから送ってもらうか、自分で取りに行くか
前章で扱った要点は、そのまま必要です。
□ at-least-once(重複する)→ 冪等性
□ ack 期限(処理が長いと再配信される)→ 期限の延長か、処理の分割
□ DLQ の設定と監視
□ 順序が必要なら順序指定キー(ただしスループットは落ちる)
どれを選ぶか
□ 仕事を配りたい、届いたら消えてよい → ブローカー型(RabbitMQ / SQS / Pub/Sub)
□ 読み直したい、複数が独立に読む → Kafka / Pub/Sub
□ 大量のイベントストリーム、順序が重要 → Kafka
□ 運用の手間をかけたくない → マネージド
□ 「この時刻にこれを実行」だけでよい → Cloud Tasks / スケジューラ(バッチとジョブ)
Kafka は強力ですが、運用コストが高いです。
□ クラスタの構成と監視
□ パーティション設計(後から減らせない)
□ コンシューマグループの管理
□ ディスクと保持期間の設計
「1日数千件のジョブを処理したい」だけなら、 シンプルなキューかマネージドサービスで十分です。
設計の基礎の「早すぎる抽象化」と同じで、 必要になってから移行するほうが安いことがほとんどです。
運用で見るもの
どの製品でも、見る指標は同じです(前章)。
□ 未処理メッセージの件数(滞留)
□ **最も古い未処理メッセージの経過時間** ← 最重要
□ DLQ の件数
□ コンシューマの処理速度とエラー率
□ Kafka なら「コンシューマラグ」(どれだけ遅れているか)
「静かに遅れる」のが非同期の怖さです。時間で監視してください。
1つの「注文確定」イベントを、在庫・通知・分析の3サービスがそれぞれ独立に処理したい。RabbitMQ ではどう構成しますか。
実務の落とし穴まとめ
- 自動 ack にする — 処理前に消えるので、落ちたら失われる
- prefetch を設定しない — 1台に偏り、増やしても速くならない
- 失敗を同じキューに戻す — 無限ループ。上限と DLQ を設ける
- DLQ を監視しない — 仕事が静かに失われる
- Kafka のパーティション数 < コンシューマ数 — 余った分は働かない
- パーティション数を後から増やす — キーの割り当てが変わり、順序に影響
- 最初から Kafka を選ぶ — 運用コストが高い。必要になってから
- コンシューマラグを見ていない — 遅れに気づけない
まとめ
- 「消えるか、残るか」が最大の違い。ここから他の性質がほぼ決まる
- RabbitMQ は仕事を配る。Exchange → Binding → Queue、ack で初めて消える。 prefetch と DLX が運用の要点
- Kafka は追記ログ。コンシューマがオフセットを持ち、読み直せる。 パーティション数が並列度の上限で、後から減らせない
- 順序が必要な単位をキーにする(Kafka)
- マネージド(Cloud Pub/Sub など)は、前章の原則がそのまま必要
- 最初から Kafka を選ばない。読み直しが要るかで判断する
- 監視するのは滞留・最古メッセージの経過時間・DLQ・ラグ
公式ドキュメント
新しい仕組みを使う時は、まず公式を1周してください(新しい言語をどう学ぶか)。
| 対象 | リンク |
|---|---|
| RabbitMQ チュートリアル | https://www.rabbitmq.com/tutorials |
| RabbitMQ ドキュメント | https://www.rabbitmq.com/docs |
| Apache Kafka 公式 | https://kafka.apache.org/documentation/ |
| Kafka 入門(Quickstart) | https://kafka.apache.org/quickstart |
| Google Cloud Pub/Sub | https://cloud.google.com/pubsub/docs |
| Amazon SQS | https://docs.aws.amazon.com/sqs/ |
RabbitMQ のチュートリアル(6本)は、1〜2時間で一通り動かせます。 概念を理解する費用対効果が非常に高いので、最初にこれをおすすめします。
章末問題
Kafka でコンシューマを5台に増やしたのに、処理速度が3台の時と変わりません。原因として最も可能性が高いのは?
次の章では、決まった時刻にまとめて動かす処理——バッチを扱います。