# Apache Kafka LinkedIn で開発され、当初はログ処理に使われた分散メッセージブローカ。Kreps・Narkhede・Rao らによる NetDB 2011 論文が原典(Kafka: A distributed messaging system for log processing)。**信頼性をスループットでトレードオフ**する設計選択を取り、高スループット・低レイテンシのメッセージキューを提供する。(Source: [[@2017__arXiv__A Survey of Distributed Message Broker Queues]]) ## アーキテクチャ(2017 時点) - データを **topic** に分割、topic を **partition** に分割し各 broker は 1 つ以上の partition を保持。consumer は **pull モデル**で自ペースで読む。 - ストレージは partition の segment 列、メッセージは**論理オフセット**で識別、broker レベルキャッシュなしで **OS ページキャッシュ**に完全依存(write-through + read-ahead)。 - consumer group が論理的に 1 consumer として動き、partition がパラレリズムの最小単位、leader/master ノード概念なし。**Zookeeper** を合意サービスとして Broker/Consumer/Ownership/Offset の 4 種レジストリで状態管理。 - 配送保証は **at-least-once**、exactly-once は 2PC が必要で「やりすぎ」と判断しデデュプは応用層に委譲。partition 内は単調増加オフセットで順序保証、メッセージ単位 CRC で破損検出。 ## スループット優位の根拠(John+ 2017 §6) 1. **SendFile API** でカーネルバッファをバイパス、file channel → socket 直書き。 2. シーケンシャルディスク書き込み + OS ページキャッシュ。 3. 標準で利用可能なバッチング。 ## ベンチマーク(John+ 2017 §5) 5 ノード(12 コア・16GB・1Gbps HDD)で Flotilla 計測: - single P/C スケールアウト: payload 一定でノード追加するとレイテンシ 3 倍改善、スループット 1.06 倍低下(ほぼ維持)。 - multi P/C 同一ノード: producer/consumer 5 倍でレイテンシ 100 倍悪化、consumer スループット 29 倍低下。Zookeeper が単一ノード上で resource contention の焦点となり context switching が増えるため、と論文は説明。 ## ログベースメッセージブローカとしてのアーキテクチャ(DDIA 2E 第12章) Kafka は*ログベースメッセージブローカ*(→ [[分散メッセージブローカ]])の代表例として、AMQP/JMS 型ブローカとは異なる設計を取る。トピックを*パーティション(partition)*(シャード)に分割し、各パーティション内でメッセージへ単調増加する*オフセット(offset)*を割り当てる。パーティション内は追記専用ログのため全順序を持つが、パーティション間の順序保証はない。コンシューマは配送されたメッセージを確認応答のたびに削除させるのではなく、定期的にオフセットを記録するだけでよく、これによりブローカ側の記帳コストを削減し高スループットを実現する——単一データベースリーダーのレプリケーションにおける log sequence number と本質的に同じ原理である。コンシューマグループは各パーティションをグループ内の1ノードへ丸ごと割り当てる粗粒度のロードバランシングを行う。(Source: [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 12 Stream Processing]] "Using logs for message storage", "Consumer offsets") ディスク上のログはセグメント単位で分割され古いセグメントは削除・アーカイブされるため、実質的に大きな固定サイズのバッファ(サーキュラーバッファ)として機能する。Kafka や Redpanda は古いメッセージをオブジェクトストレージへの階層化ストレージとして退避する機能も持ち、WarpStream・Confluent Freight・Bufstream のように全データをオブジェクトストレージに置く構成(メッセージを Iceberg テーブルとして格納しバッチ・データウェアハウスジョブから直接参照可能にする)も登場している。(Source: [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 12 Stream Processing]] "Disk space usage") *ログコンパクション(log compaction)*機能により、各キーの最新値のみを残しつつ他のキーは無期限に保持できるため、Kafka を単なる一時的メッセージングでなく永続ストレージとして使うことができる。CDC(→ [[変更データキャプチャ(CDC)]])のイベントストリームにログコンパクションを適用すると、オフセット0からの走査だけでデータベースの全内容を再構築でき、別途スナップショットを取得する必要がなくなる。Kafka Streams と Confluent の ksqlDB はこのログコンパクションを前提としたマテリアライズドビュー保守機能を提供する。(Source: [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 12 Stream Processing]] "Log compaction", "Maintaining materialized views") Kafka は Google Cloud Dataflow・VoltDB と並び、ストリーム処理における厳密に一度(exactly-once、実際には effectively-once)のセマンティクスをアトミックコミットで実現する実装の一つでもある。XA のような異種システム間のトランザクションとは異なり、状態変更とメッセージングを Kafka 自身の内部に限定することで効率的なアトミックコミットを実現する(トランザクショナルメッセージング)。Kafka Streams はさらに、オペレータの状態をログコンパクション付きの専用 Kafka トピックへ変更送信することで障害後の状態複製を実現する。(Source: [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 12 Stream Processing]] "Atomic commit revisited", "Rebuilding state after a failure") ## データ統合・アンバンドリングの標準プロトコルとしての位置づけ(DDIA 2E 第13章) 第13章はKafkaを個別のアーキテクチャ詳細としてではなく、データベースのアンバンドリング(→ [[データベースのアンバンドリング]])を実務上可能にしている基盤技術の1つとして位置づける。Debezium(→ [[Debezium]])が多様なデータベースから変更ストリームを抽出し、Kafkaのプロトコルがイベントストリームの事実上の標準になりつつあることで、異種データシステム間の統合(→ [[データ統合]])のコストが下がりつつあると述べる。また、イベントソーシング(→ [[イベントソーシングとCQRS]])の実装においてイベントログとして使われるメッセージブローカの代表例としても言及される。(Source: [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 13 A Philosophy of Streaming Systems]] "Making unbundling work") ## 関連 - ソース: [[@2017__arXiv__A Survey of Distributed Message Broker Queues]] / [[@2026__Netflix TechBlog__From Silos to Service Topology - Why Netflix Built a Real-Time Service Map]] / [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 12 Stream Processing]] / [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 13 A Philosophy of Streaming Systems]] - 概念: [[分散メッセージブローカ]] / [[サービストポロジ]] / [[変更データキャプチャ(CDC)]] / [[ストリーム処理の耐障害性]] / [[データ統合]] / [[データベースのアンバンドリング]] / [[イベントソーシングとCQRS]] - 比較対象: [[AMQP]]、[[RabbitMQ]] - 信頼性テスト: [[Sieve]] / [[フェイルスローハードウェア]] - 隣接製品: [[Debezium]](Kafka Connect 経由でしばしば併用される CDC プラットフォーム) ## wiki 内の言及 - [[@2026__Netflix TechBlog__From Silos to Service Topology - Why Netflix Built a Real-Time Service Map]]: Netflix Service Topology の多リージョン取り込みレイヤーとして採用。各 AWS リージョンから eBPF フローログを消費する。 - `Observability Engineering`(2nd Edition)第7章: [[OpenTelemetry]] のストリーミングアーキテクチャ計装例として、Kafka のようなメッセージブローカを介したイベント駆動パイプラインを取り上げる。publish 側はメッセージ envelope へトレースコンテキストをシリアライズし(`messaging.publish` producer スパン)、consume 側は envelope からコンテキストを抽出して `messaging.process` スパンを生成、両者をスパンリンクで(直接の親子関係ではなく)結び付ける設計を推奨する。(Source: [[@2026__OReilly__Observability Engineering 2E - Chapter 7 Instrumenting Your Code with OpenTelemetry]]) - `Designing Data-Intensive Applications`(2nd Edition)第12章: ログベースメッセージブローカの代表例として、パーティション/オフセットモデル・ログコンパクション・階層化ストレージ・トランザクショナルメッセージングによる exactly-once を詳述。(Source: [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 12 Stream Processing]]) - `Designing Data-Intensive Applications`(2nd Edition)第13章: データベースのアンバンドリングを支える「事実上の標準プロトコル」として言及。Debezium と組み合わせた CDC パイプライン、イベントソーシングのイベントログとしての利用に触れる。(Source: [[@2026__OReilly__Designing Data-Intensive Applications 2E - Chapter 13 A Philosophy of Streaming Systems]]) - [[@2026__TOCS__Fail-Slow Hardware Failure Bug Analysis and Detection for Cloud Systems]]: フェイルスローハードウェア障害注入テスト [[Sieve]] の評価対象(Kafka 3.6.0、producer/consumer 性能テストをワークロードとする)。未知バグ 2 件を検出した——KAFKA-16401 は consumer 要求が `storeOffsets` で詰まりグループロックを保持したままになり、タイムアウトによる再送で全 request handler スレッドが占有される障害。KAFKA-16412 は同一トピックの並行作成要求で 2 番目が `TopicExistsException` を返すが、未確立のトピックに対する後続要求が失敗し 2 番目のクライアントを混乱させる障害(開発者に確認済み)。Kafka は同論文のバグスタディ対象 5 システムには含まれておらず、手法が研究対象システム固有でないことを示すベンチマークとして選ばれた。 - [[@2025__HOTOS__Understanding the limitations of pubsub systems]]: pubsub システムの限界を論じる批判的立場から Kafka を代表例として扱う。トピックコンパクションのサポートを挙げつつ、通知なしにコンパクションが起きるため subscriber が未見イベントの消失に気づけない点、改良版 Kafka が API を変更して暗黙のストレージ層を明示化すれば著者らの提案する「ingestion storage with built-in watch」象限に収まると指摘する。