第6章 Kafka クライアントの開発


任意のプログラミング言語で Kafka クライアントを作成し、AMQ Streams に接続します。

Kafka クラスターとやりとりするには、クライアントアプリケーションがメッセージを生成および消費できる必要があります。基本的な Kafka クライアントアプリケーションを開発して設定するには、少なくとも次のことを行う必要があります。

  • Kafka クラスターに接続するための設定をセットアップする
  • プロデューサーとコンシューマーを使用してメッセージを送受信する

Kafka クラスターに接続し、プロデューサーとコンシューマーを使用するための基本設定をセットアップすることは、Kafka クライアント開発の最初のステップです。その後、入力、セキュリティー、パフォーマンス、エラー処理、クライアントアプリケーションの機能の改善に拡張できます。

前提条件

以下のプロパティー値を含むクライアントプロパティーファイルが作成されました。

手順

  1. Java、Python、.NET などのプログラミング言語用の Kafka クライアントライブラリーを選択します。
  2. パッケージマネージャーを使用するか、ソースからライブラリーをダウンロードして手動でライブラリーをインストールします。
  3. Kafka クライアントに必要なクラスと依存関係をコードにインポートします。
  4. 作成するクライアントのタイプに応じて、Kafka コンシューマーオブジェクトまたはプロデューサーオブジェクトを作成します。

    両方を実行するクライアントを使用することもできます。

  5. Kafka クラスターに接続するための設定プロパティー (必要に応じてブローカーアドレス、ポート、認証情報など) を指定します。
  6. Kafka コンシューマーまたはプロデューサーオブジェクトを使用して、トピックのサブスクライブ、メッセージの生成、または Kafka クラスターからのメッセージの取得を行います。
  7. AMQ Streams との接続または通信中に発生する可能性のあるエラーを処理します。

6.1. Kafka プロデューサークライアントの例

この Java ベースの Kafka プロデューサークライアントは、Kafka トピックへのメッセージを生成する自己完結型アプリケーションの例です。クライアントは Kafka Producer API を使用して、いくつかのエラー処理を行いながらメッセージを非同期に送信します。

クライアントは、メッセージ処理用の Callback インターフェイスを実装します。

Kafka プロデューサークライアントを実行するには、Producer クラスの main メソッドを実行します。クライアントは、randomBytes メソッドを使用して、メッセージペイロードとしてランダムなバイト配列を生成します。クライアントは、NUM_MESSAGES 個のメッセージ (設定例では 100) が送信されるまで、Kafka トピックへのメッセージを生成します。プロデューサはスレッドセーフであるため、複数のスレッドが単一のプロデューサインスタンスを使用できます。

このサンプルクライアントは、特定のユースケース向けに、より複雑な Kafka プロデューサを構築するための基本基盤を提供します。ロギングフレームワークとの統合など、追加の機能を組み込むことができます。

注記

SLF4J バインディングを各クライアントに追加して、クライアント API ログを表示できます。

前提条件

  • 指定された BOOTSTRAP_SERVERS で実行されている Kafka ブローカー
  • メッセージが生成される TOPIC_NAME という名前の Kafka トピック。

設定

プロデューサークライアントは、Producer クラスで指定された以下の定数で設定できます。

BOOTSTRAP_SERVERS
Kafka ブローカーに接続するためのアドレスとポート (例: localhost:9092)。
TOPIC_NAME
メッセージを生成する Kafka トピックの名前。
NUM_MESSAGES
停止する前に生成するメッセージの数。
MESSAGE_SIZE_BYTES
各メッセージのバイト単位のサイズ。
PROCESSING_DELAY_MS
メッセージ送信間の遅延 (ミリ秒単位)。これにより、メッセージの処理時間をシミュレートでき、テストに役立ちます。

プロデューサークライアントの例

import java.util.Properties;
import java.util.Random;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;

import org.apache.kafka.clients.producer.Callback;
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.clients.producer.RecordMetadata;
import org.apache.kafka.common.errors.RetriableException;
import org.apache.kafka.common.serialization.ByteArraySerializer;
import org.apache.kafka.common.serialization.LongSerializer;

public class Producer implements Callback {
    private static final Random RND = new Random(0);
    private static final String BOOTSTRAP_SERVERS = "localhost:9092";
    private static final String TOPIC_NAME = "my-topic";
    private static final long NUM_MESSAGES = 100;
    private static final int MESSAGE_SIZE_BYTES = 100;
    private static final long PROCESSING_DELAY_MS = 0L;

    protected AtomicLong messageCount = new AtomicLong(0);

    public static void main(String[] args) {
        new Producer().run();
    }

    public void run() {
        System.out.println("Running producer");
        try (var producer = createKafkaProducer()) {  
1

            byte[] value = randomBytes(MESSAGE_SIZE_BYTES); 
2

            while (messageCount.get() < NUM_MESSAGES) { 
3

                sleep(PROCESSING_DELAY_MS); 
4

                producer.send(new ProducerRecord<>(TOPIC_NAME, messageCount.get(), value), this); 
5

                messageCount.incrementAndGet();
            }
        }
    }

    private KafkaProducer<Long, byte[]> createKafkaProducer() {
        Properties props = new Properties(); 
6

        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS); 
7

        props.put(ProducerConfig.CLIENT_ID_CONFIG, "client-" + UUID.randomUUID()); 
8

        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class); 
9

        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class);
        return new KafkaProducer<>(props);
    }

    private void sleep(long ms) { 
10

        try {
            TimeUnit.MILLISECONDS.sleep(ms);
        } catch (InterruptedException e) {
            throw new RuntimeException(e);
        }
    }

    private byte[] randomBytes(int size) { 
11

        if (size <= 0) {
            throw new IllegalArgumentException("Record size must be greater than zero");
        }
        byte[] payload = new byte[size];
        for (int i = 0; i < payload.length; ++i) {
            payload[i] = (byte) (RND.nextInt(26) + 65);
        }
        return payload;
    }

    private boolean retriable(Exception e) { 
12

        if (e == null) {
            return false;
        } else if (e instanceof IllegalArgumentException
            || e instanceof UnsupportedOperationException
            || !(e instanceof RetriableException)) {
            return false;
        } else {
            return true;
        }
    }

    @Override
    public void onCompletion(RecordMetadata metadata, Exception e) { 
13

        if (e != null) {
            System.err.println(e.getMessage());
            if (!retriable(e)) {
                e.printStackTrace();
                System.exit(1);
            }
        } else {
            System.out.printf("Record sent to %s-%d with offset %d%n",
                metadata.topic(), metadata.partition(), metadata.offset());
        }
    }
}

1
クライアントは、createKafkaProducer メソッドを使用して Kafka プロデューサを作成します。プロデューサは、Kafka トピックにメッセージを非同期的に送信します。
2
バイト配列は、Kafka トピックに送信される各メッセージのペイロードとして使用されます。
3
送信されるメッセージの最大数は、NUM_MESSAGES 定数値によって決まります。
4
メッセージレートは、送信される各メッセージ間の遅延によって制御されます。
5
プロデューサは、トピック名、メッセージカウント値、およびメッセージ値を渡します。
6
クライアントは、提供された設定を使用して KafkaProducer インスタンスを作成します。プロパティーファイルを使用することも、設定を直接追加することもできます。基本設定の詳細については、4章Kafka クラスターに接続するためのクライアントアプリケーションの設定 を参照してください。
7
Kafka ブローカーへの接続。
8
ランダムに生成された UUID を使用するプロデューサーの一意のクライアント ID。クライアント ID は必須ではありませんが、リクエストのソースを追跡するのに役立ちます。
9
キーおよび値をバイト配列として処理するための適切なシリアライザークラス。
10
指定されたミリ秒数の間、メッセージ送信プロセスに遅延を導入するメソッド。メッセージの送信を担当するスレッドが一時停止中に中断されると、InterruptedException エラーがスローされます。
11
Kafka トピックに送信される各メッセージのペイロードとして機能する、特定のサイズのランダムなバイト配列を作成するメソッド。このメソッドはランダムな整数を生成し、ASCII コードの大文字を表す 65 を加算します (65 は A、66 は B など)。ASCII コードはペイロード配列に 1 バイトとして保存されます。ペイロードサイズがゼロ以下の場合、IllegalArgumentException がスローされます。
12
例外の後にメッセージの送信を再試行するかどうかを確認するメソッド。Null 例外と指定された例外は再試行されません。また、RetriableException インターフェイスを実装していない例外も再試行されません。このメソッドをカスタマイズして、他のエラーを含めることができます。
13
Kafka ブローカーによってメッセージが確認されたときに呼び出されるメソッド。成功すると、トピック、パーティション、メッセージのオフセット位置の詳細を含むメッセージが出力されます。メッセージの送信時にエラーが発生した場合は、エラーメッセージが出力されます。このメソッドは例外をチェックし、それが致命的なエラーであるか致命的ではないエラーであるかに基づいて適切なアクションを実行します。エラーが致命的ではない場合、メッセージ送信プロセスは続行されます。エラーが致命的である場合、スタックトレースが出力され、プロデューサは終了します。

エラー処理

プロデューサークライアントがキャッチした致命的な例外:

InterruptedException
一時停止中に現在のスレッドが中断された場合にスローされるエラー。通常、中断はプロデューサーを停止またはシャットダウンするときに発生します。例外は RuntimeException として再スローされ、プロデューサーが終了します。
IllegalArgumentException
プロデューサが無効または不適切な引数を受け取ったときにスローされるエラー。たとえば、トピックが欠落している場合、例外がスローされます。
UnsupportedOperationException
操作がサポートされていない場合、またはメソッドが実装されていない場合にスローされるエラー。たとえば、サポートされていないプロデューサー設定を使用しようとしたり、KafkaProducer クラスでサポートされていないメソッドを呼び出そうとした場合、例外がスローされます。

プロデューサークライアントによってキャッチされた致命的ではない例外:

RetriableException
Kafka クライアントライブラリーによって提供される RetriableException インターフェイスを実装する例外に対してスローされるエラー。

致命的ではないエラーの場合、プロデューサーはメッセージの送信を続けます。

Red Hat logoGithubredditYoutubeTwitter

詳細情報

試用、購入および販売

コミュニティー

会社概要

Red Hat は、企業がコアとなるデータセンターからネットワークエッジに至るまで、各種プラットフォームや環境全体で作業を簡素化できるように、強化されたソリューションを提供しています。

多様性を受け入れるオープンソースの強化

Red Hat では、コード、ドキュメント、Web プロパティーにおける配慮に欠ける用語の置き換えに取り組んでいます。このような変更は、段階的に実施される予定です。詳細情報: Red Hat ブログ.

Red Hat ドキュメントについて

Legal Notice

Theme

© 2026 Red Hat
トップに戻る