プログラマのための IT 教科書

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(受け取った瞬間に消える)は、失うことを許容できる場合だけ
prefetch を設定しないと1台に偏る

既定では、コンシューマが大量のメッセージをまとめて抱え込むことがあります。

コンシューマ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      各パーティション内の通し番号

決定的に違う点

RabbitMQKafka
読んだら消える消えない(保持期間まで残る)
進捗の管理ブローカーが持つコンシューマがオフセットを持つ
読み直しできないできる(オフセットを戻す)
順序保証しにくいパーティション内では保証される
複数の購読者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 SNSfanout通知の配信
Cloud Tasksタスク実行「この URL をこの時刻に叩く」
Amazon MSK / ConfluentKafkaKafka そのものをマネージドで
Cloud Pub/Sub の考え方

Google Cloud を使うなら、まずこれが選択肢になります。

Topic          送信先
Subscription   購読の単位。**トピックごとに複数作れる**
               → サブスクリプションを分ければファンアウトになる
Push / Pull    サーバーから送ってもらうか、自分で取りに行くか

前章で扱った要点は、そのまま必要です。

□ at-least-once(重複する)→ 冪等性
□ ack 期限(処理が長いと再配信される)→ 期限の延長か、処理の分割
□ DLQ の設定と監視
□ 順序が必要なら順序指定キー(ただしスループットは落ちる)

どれを選ぶか

□ 仕事を配りたい、届いたら消えてよい     → ブローカー型(RabbitMQ / SQS / Pub/Sub)
□ 読み直したい、複数が独立に読む         → Kafka / Pub/Sub
□ 大量のイベントストリーム、順序が重要   → Kafka
□ 運用の手間をかけたくない               → マネージド
□ 「この時刻にこれを実行」だけでよい     → Cloud Tasks / スケジューラ(バッチとジョブ)
最初から Kafka を選ばない

Kafka は強力ですが、運用コストが高いです。

□ クラスタの構成と監視
□ パーティション設計(後から減らせない)
□ コンシューマグループの管理
□ ディスクと保持期間の設計

「1日数千件のジョブを処理したい」だけなら、 シンプルなキューかマネージドサービスで十分です。

設計の基礎の「早すぎる抽象化」と同じで、 必要になってから移行するほうが安いことがほとんどです。

運用で見るもの

どの製品でも、見る指標は同じです(前章)。

□ 未処理メッセージの件数(滞留)
□ **最も古い未処理メッセージの経過時間**  ← 最重要
□ DLQ の件数
□ コンシューマの処理速度とエラー率
□ Kafka なら「コンシューマラグ」(どれだけ遅れているか)

「静かに遅れる」のが非同期の怖さです。時間で監視してください。

1つの「注文確定」イベントを、在庫・通知・分析の3サービスがそれぞれ独立に処理したい。RabbitMQ ではどう構成しますか。

実務の落とし穴まとめ

  1. 自動 ack にする — 処理前に消えるので、落ちたら失われる
  2. prefetch を設定しない — 1台に偏り、増やしても速くならない
  3. 失敗を同じキューに戻す — 無限ループ。上限と DLQ を設ける
  4. DLQ を監視しない — 仕事が静かに失われる
  5. Kafka のパーティション数 < コンシューマ数 — 余った分は働かない
  6. パーティション数を後から増やす — キーの割り当てが変わり、順序に影響
  7. 最初から Kafka を選ぶ — 運用コストが高い。必要になってから
  8. コンシューマラグを見ていない — 遅れに気づけない

まとめ

  • 「消えるか、残るか」が最大の違い。ここから他の性質がほぼ決まる
  • 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/Subhttps://cloud.google.com/pubsub/docs
Amazon SQShttps://docs.aws.amazon.com/sqs/

RabbitMQ のチュートリアル(6本)は、1〜2時間で一通り動かせます。 概念を理解する費用対効果が非常に高いので、最初にこれをおすすめします。

章末問題

Kafka でコンシューマを5台に増やしたのに、処理速度が3台の時と変わりません。原因として最も可能性が高いのは?

次の章では、決まった時刻にまとめて動かす処理——バッチを扱います。

読み終わったら記録しておくと、目次で進み具合が分かります。