【かふかすとりーむず】

Kafka Streams とは?

最終更新:
💡 Kafkaの流れを、アプリに組み込んだ処理で加工する

Kafkaのデータを継続的に変換・集計する処理を、JVMアプリに組み込むクライアントライブラリ。Kafkaブローカーとは別に、処理するアプリの実行環境が必要。

📌 このページのポイント
処理は、ブローカーとは別のアプリで 🖥 JVMアプリ Kafka Streams 条件で絞る 項目を変換する 商品ごとに集計 状態ストア 🗄 Kafkaブローカー側 入力トピック 出力トピック 読み取り 書き込み 専用の処理クラスタは不要 アプリの実行環境と運用は必要
注文の変換・集計アプリの模式図。枠でアプリとKafkaブローカーを分け、矢印はトピックの読み取りと結果の書き込みを示す。状態ストアはアプリ内にあり、復元用の変更ログなどは省略。
ひよこ ひよこ
KafkaとKafka Streamsはどう違うの?
ペンギン先生 ペンギン先生
Kafkaはイベントを保存して届ける基盤で、Kafka Streamsはそのイベントを加工する処理をアプリに組み込むライブラリだよ。ブローカー自体がアプリの集計コードを動かす、という意味ではないんだ。
ひよこ ひよこ
どんな加工ができるの?
ペンギン先生 ペンギン先生
注文を条件で絞る、必要な項目に変換する、商品ごとの金額を集計する、といった処理だよ。Kafkaの入力トピックから読み、結果を出力トピックへ送るアプリを作れるんだ。
ひよこ ひよこ
専用クラスタが不要なら、サーバーもいらないの?
ペンギン先生 ペンギン先生
処理するアプリを動かすマシンやコンテナなどは必要だよ。Kafka Streams専用の処理クラスタを用意する必要がない、という意味なんだ。Kafkaがあってもアプリの資源や配置、監視まで不要にはならないよ。
ひよこ ひよこ
集計の途中の値はどこに置くの?
ペンギン先生 ペンギン先生
ローカルの状態ストアに保持するよ。永続ストアにはRocksDBを使う構成があり、変更ログをKafkaのトピックに記録して復元に使う。定期的な丸ごとバックアップとは違い、変更ログの設定や保持も関係するんだ。
ひよこ ひよこ
アプリを増やせばいくらでも速くなる?
ペンギン先生 ペンギン先生
処理をタスクに分け、複数のアプリインスタンスへ割り当てるよ。入力トピックのパーティション数などが並列化の上限に関係する。アプリを増やすだけで無制限に速くなるわけではないんだ。
ペンギン
まとめ:ざっくりこれだけ覚えればOK!
「Kafka Streams」って出てきたら「Kafkaのデータ処理をJVMアプリに組み込むライブラリ」と思えばだいたいOK!
📖 おまけ:英語の意味
「Kafka Streams」 = Kafkaのストリーム処理ライブラリ
💬 Apache Kafkaのサブプロジェクトとして開発されたライブラリだよ。streams(流れ)という名の通り、データを流れとして連続処理するんだ

参考資料

← 用語集にもどる