Kafkaの登場背景とメッセージングシステムの利点
Apache Kafkaは、元々LinkedInにおいてアクティビティストリームや運用データのパイプライン処理を目的として開発された分散メッセージングシステムです。PV(ページビュー)や検索クエリなどのユーザー行動データ、およびCPU使用率やIOなどのサーバーメトリクスを効率的に処理するために設計されました。
メッセージングシステムを導入する主なメリットは以下の通りです:
- デカップリング:データに基づくインターフェース層を介在させることで、システム間の結合度を下げます。
- データの永続化と冗長性:メッセージをキューに保持し、データ損失を防ぎます。
- スケーラビリティ:処理プロセスを分離することで、システムの拡張が容易になります。
- フォールトトレランス:コンシューマーがダウンしても、復旧後に未処理のメッセージから再開できます。
- 順序の保証:同一パーティション内でのメッセージ順序を厳密に保証します。
- 非同期通信:メッセージをキューに蓄積し、必要に応じて非同期に処理できます。
コアコンセプトとアーキテクチャ
Kafkaのアーキテクチャは、プロデューサー、ブローカー、コンシューマーグループ、そしてZooKeeperで構成されます。プロデューサーはメッセージをブローカーにプッシュし、コンシューマーはブローカーからメッセージをプルします。
主要な用語は以下の通りです:
- Broker:Kafkaクラスタを構成するサーバーインスタンス。
- Topic:メッセージを分類するための論理的なチャネル。
- Partition:トピックを物理的に分割した単位。並列処理とスケーラビリティの基盤となります。
- Producer / Consumer:メッセージの送信者および受信者。
- Consumer Group:協調してメッセージを消費するコンシューマーの集合体。
主なユースケース
- リアルタイムストリーム処理:システム間でデータを確実に転送するパイプラインの構築。
- 行動トラッキング:ユーザーのクリックや検索行動をリアルタイムに記録し、ビッグデータプラットフォームで分析。
- ログ集約:分散サーバーのログをKafkaに集約し、ElasticsearchやHDFSに転送して検索やオフライン分析に活用。
環境構築と基本操作
1. インストールと起動
公式サイトからアーカイブをダウンロードし、展開します。
tar -xzf kafka_2.13-3.4.0.tgz
cd kafka_2.13-3.4.0
KafkaはZooKeeperに依存しているため、まずZooKeeperを起動し、続いてKafkaブローカーを起動します。
# ZooKeeperの起動
bin/zookeeper-server-start.sh config/zookeeper.properties &
# Kafkaブローカーの起動
bin/kafka-server-start.sh config/server.properties &
2. トピックの管理とメッセージの送受信
トピックの作成、一覧表示、およびコンソールからの送受信は以下のコマンドで行います。
# トピックの作成
bin/kafka-topics.sh --create --topic user-events --bootstrap-server localhost:9092 --partitions 3 --replication-factor 1
# トピック一覧の表示
bin/kafka-topics.sh --list --bootstrap-server localhost:9092
# メッセージの送信
bin/kafka-console-producer.sh --topic user-events --bootstrap-server localhost:9092
# メッセージの受信
bin/kafka-console-consumer.sh --topic user-events --from-beginning --bootstrap-server localhost:9092
3. クラスタ構成の設定
複数ノードでクラスタを構築する場合、各ノードのserver.propertiesで以下の設定を調整します。
broker.id:各ブローカーに一意のID(0, 1, 2...)を設定。advertised.listeners:クライアントが接続するための外部IPアドレスとポートを指定(例:PLAINTEXT://192.168.1.10:9092)。
Javaクライアントによる実装
Maven依存関係の追加
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.4.0</version>
</dependency>
プロデューサーの実装
以下のコードは、非同期でメッセージを送信するプロデューサーの例です。設定はMapを利用して構築し、クリーンな構造にしています。
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.HashMap;
import java.util.Map;
public class EventProducer {
public static void main(String[] args) {
Map<String, Object> config = new HashMap<>();
config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
config.put(ProducerConfig.ACKS_CONFIG, "all");
String targetTopic = "user-events";
try (KafkaProducer<String, String> producer = new KafkaProducer<>(config)) {
for (int i = 0; i < 100; i++) {
String payload = "event_data_" + i;
ProducerRecord<String, String> record = new ProducerRecord<>(targetTopic, String.valueOf(i), payload);
producer.send(record, (metadata, exception) -> {
if (exception == null) {
System.out.printf("Sent to partition %d, offset %d%n", metadata.partition(), metadata.offset());
} else {
exception.printStackTrace();
}
});
}
}
}
}
コンシューマーの実装
コンシューマーグループに参加し、 지속적으로メッセージをポーリングして処理する例です。
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
public class EventConsumer {
public static void main(String[] args) {
Map<String, Object> config = new HashMap<>();
config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
config.put(ConsumerConfig.GROUP_ID_CONFIG, "event-processing-group");
config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
config.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(config)) {
consumer.subscribe(Collections.singletonList("user-events"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("Consumed message: key=%s, value=%s, partition=%d%n",
record.key(), record.value(), record.partition());
}
}
}
}
}
プロデューサーの重要設定パラメータ
- acks:メッセージ送信時の確認応答レベル。
0:応答を待たない(低レイテンシ、データ損失のリスクあり)。1:リーダーレプリカからの応答のみ待機。all(または-1):すべてのISR(In-Sync Replicas)からの応答を待機(高耐久性)。
- batch.size:同一パーティションへのメッセージをバッチ送信する際のバッファサイズ(デフォルト16KB)。
- linger.ms:バッチ送信を遅延させる時間。設定することで、より多くのメッセージをまとめて送信し、スループットを向上させます。
- max.request.size:1回のリクエストで送信できる最大メッセージサイズ(デフォルト1MB)。