なぜ Kafka を触ろうと思ったか
きっかけは、同じ場面に何度も出くわしたことだった。
何かのイベントが起きたときに、それを別のサービスにも渡したくなる。注文が入ったら在庫を減らして、通知も送って、あとから「分析にも使いたい」と言われる。そのたびに呼び出し先を1つずつ増やしていくと、送る側が受け取る側を全部知っている状態ができあがる。受け取る側が1つ落ちると送信側まで巻き込まれるし、あとから1つ足すのに送信側を触らないといけない。
Kafka の名前はずっと聞いていた。ただ、なんとなく「大量のトラフィックを捌く大きな会社が使うもの」だと思っていて、自分の手元で触るものだと思っていなかった。
実際は、Docker があればローカルにブローカー1台立てて全部試せる。それを知って、7つのユースケースを順番に動かす検証用のリポジトリを作った。この記事は、そのときに「思っていたのと違った」ことの記録だ。
最初は「すごいメッセージキュー」だと思っていた
触る前のイメージはこうだった。キューにメッセージを入れる。受け取る側が取り出すと、そのメッセージは消える。複数のワーカーで並列に捌ける。要するに RabbitMQ や SQS の高性能版だろう、と。
これが根本から違った。
Kafka のトピックはキューではなく、追記専用のログだ。書き込まれたメッセージは読まれても消えない。保存期間(デフォルトは7日)が来るまでそこに残り続ける。
では読み手はどう管理されるのか。各メッセージには 0 から始まる連番(オフセット)が振られていて、読み手側が「どこまで読んだか」を記録するだけだ。メッセージを削除するのではなく、しおりを進めていくイメージに近い。
offset: 0 1 2 3 4
[msg A] [msg B] [msg C] [msg D] [msg E]
↑
ここまで読んだ、という記録だけを持つ
この違いが効いてくるのは、同じデータを別の用途で最初から読み直せるところだ。分析チームがあとから参加して、過去7日ぶんを頭から読み直せる。「取り出したら消える」前提で設計を考えていると、この嬉しさがまったく理解できない。自分は最初そこで詰まっていた。
つまずいたのは、何も表示されないこと
概念の話より先に、実際の手触りの話をする。
環境を立ち上げて、サンプルのコンシューマーを起動した。何も表示されない。
python use_cases/01_basic_pubsub/consumer.py
# ...無言
これ、正常動作だ。コンシューマーはメッセージが来るのを待ち続けるループなので、誰も送っていなければ当然何も出ない。別のターミナルを開いてプロデューサーを動かすと、そこで初めて流れ始める。
いま思えば当たり前なのだが、キューの発想だと「取りに行ったら何か返ってくる」と思ってしまう。ターミナルを2枚開いて初めて動く、というのが最初の壁だった。
もう1つ。先に流したメッセージが読めない。 プロデューサーを動かしてからコンシューマーを起動すると、やっぱり何も出ない。
これは auto.offset.reset のせいで、デフォルトが latest だからだ。つまり「起動した時点より後に来たものだけ読む」。過去に遡って読みたいなら earliest にする。自分のリポジトリでは、テンプレートの時点で earliest を既定にした。学習中は過去が読めないと何が起きたのか分からないので。
この2つを踏んだところで、**「これはキューじゃなくてログなんだ」**というのが頭ではなく手で分かった。ログだから、読み始める位置を自分で決められる。読み終わっても残っている。
ちなみに「コンシューマーの止め方が分からない」も通った。あれは無限ループなので Ctrl+C でいい。
順序が保証されるのは、トピック単位ではない
「Kafka は順序を保証する」という説明はよく見る。これも、そのまま受け取ると事故る。
トピックは内部で複数のパーティションに分かれる。分けるほど並列に捌けるようになる。ただし順序が保証されるのはパーティションの中だけで、トピック全体で見た順序は保証されない。3つに分けた時点で、全体の順番という概念が無くなる。
ここで効いてくるのがキーだ。メッセージにキーを付けて送ると、同じキーのものは必ず同じパーティションに入る。
Topic: orders(パーティション 3 つ)
Partition 0: [order_A の作成] [order_A の支払い] [order_A の発送]
Partition 1: [order_B の作成] [order_B の支払い]
Partition 2: [order_C の作成]
注文IDをキーにすれば、同じ注文に関するイベントは必ず順番どおりに処理される。違う注文どうしの前後関係は保証されないけれど、そもそも保証する必要がない。
つまりキー設計が本体だった。「順序を守りたい単位は何か」を決めてから、それをキーにする。ここを決めずにパーティションだけ増やすと、速いけど順番が狂うシステムができあがる。
Consumer Group が2つの意味を兼ねている
ここも戸惑った。
複数のコンシューマーを同じグループにまとめると、パーティションが自動的に分担される。各メッセージはグループ内で1回だけ処理される。これは想像どおり。
一方、グループが違えば、同じメッセージを全グループが受け取る。
orders ──→ [Group: inventory] 全メッセージを受信
──→ [Group: notification] 全メッセージを受信
──→ [Group: analytics] 全メッセージを受信
つまり「負荷分散」と「ファンアウト」を、グループ名という1つの設定で切り替えている。キューとトピックを別々の概念として持つミドルウェアから来ると、ここが一番混乱すると思う。
怖いのは、グループ名を1文字変えるだけで挙動がまるごと変わるのに、エラーは何も出ないことだ。分担させるつもりが全員に配られていた、という事故が起こりえる。
Exactly-Once はデフォルトではない
配信保証は3種類ある。
| 保証レベル | 重複 | 消失 | 用途 |
|---|---|---|---|
| At-Most-Once | なし | あり | 多少欠けてもいいログ |
| At-Least-Once | あり | なし | 通知・分析 |
| Exactly-Once | なし | なし | 課金・在庫 |
デフォルトは At-Least-Once。 つまり重複はありうる。「Kafka を使えば重複しない」ではない。Exactly-Once はトランザクション機能を使って初めて成立するもので、何もしなければ得られない。
なので、受け取る側を冪等に書くのは結局やる。ここは他のメッセージング基盤と変わらなかった。
ローカルで作ってから、クラウドに寄せる
環境は Docker Compose でローカルに立てた。ブローカー1台、Schema Registry、それとブラウザで中を覗ける UI。学習目的ならこれで十分だ。ブローカーが1台でも、ここまで書いた概念は全部試せる。
そのうえで、Confluent Cloud(マネージドの SaaS)へ切り替えられるように、接続設定を環境変数に外出ししておいた。コードを変えずに .env だけで移れる形だ。
移すときに知っておきたかったのは、必要なのが「接続先の変更」だけではないことだった。
KAFKA_BOOTSTRAP_SERVERS=pkc-xxxxx.<region>.aws.confluent.cloud:9092
KAFKA_SECURITY_PROTOCOL=SASL_SSL # ローカルは PLAINTEXT
KAFKA_SASL_MECHANISM=PLAIN
KAFKA_SASL_USERNAME=<クラスタの API キー>
KAFKA_SASL_PASSWORD=<クラスタの API シークレット>
SCHEMA_REGISTRY_URL=https://psrc-xxxxx.confluent.cloud
SCHEMA_REGISTRY_API_KEY=<Schema Registry の API キー> # クラスタのとは別物
SCHEMA_REGISTRY_API_SECRET=<Schema Registry の API シークレット>
ローカルは平文の接続なので、クラウドに向けると SASL_SSL に変わる。そして認証情報が2組いる。 クラスタ用の API キーと、Schema Registry 用の API キーは別物だ。ここを知らずに「キーは発行した」と思って進むと、ブローカーには繋がるのにスキーマの解決だけ落ちる、という分かりにくい詰まり方をする。自分はここで30分ほど無駄にした。
逆に言えば、ローカルで概念を掴んでからクラウドに行くほうが確実に安い。 無料枠にも限りがあるので、ループを回して壊しながら覚える段階はローカルでやったほうがいい。
これから触ってみる人へ
Kafka を触ってみようか迷っているなら、ローカルの Docker で十分だ。ブローカー1台で、ここまで書いたことは全部確認できる。
最初にやるべきなのは、「読んでも消えない」を体で確かめることだと思う。プロデューサーでメッセージを何件か流して、コンシューマーで読む。止める。そしてグループ名を変えて、もう一度起動する。同じメッセージがもう一度、頭から流れてくる。
これを見た瞬間に、キューとの違いが腹に落ちた。自分はここからようやく、パーティションとキー、Consumer Group、配信保証と順番に進めるようになった。
大きなトラフィックを捌くために使うもの、という印象が強い技術だが、実際に触ってみると**「データを消さずに、あとから誰でも読み直せる場所」**という性格のほうが本体だった。規模の話は、そのあとについてくるものだと思う。