Apache Kafkaの基礎アーキテクチャとJavaによる実践ガイド

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)。

タグ: Apache Kafka Java 分散メッセージング ストリーム処理 ZooKeeper

8月13日 22:03 投稿