Kafka 割り当て量の設定

Kafka 割り当て量は、クライアントによって使用されるブローカー・リソースを制御するために、作成要求および消費要求に制限を適用します。 Kafka 割り当て量により、管理者は、個々のプロデューサー・アプリケーションおよびコンシューマー・アプリケーションが消費できるネットワーク・スループットに制限を適用できます。

Kafka 割り当て量について

制約のないままにしておくと、少数のコンシューマーまたはプロデューサーがサービス・インスタンスの使用可能なネットワーク・スループットを独占する可能性があります。

Kafka ブローカーは、クライアントがネットワークを飽和状態にしたり、ブローカー・リソースを独占したりしないようにレート制限を適用する割り当て量をサポートします。 詳しくは、以下を参照してください。 Apache Kafka の資料。

Kafka 割り当て量は、ネットワーク帯域幅の使用量を制限するように構成できます。 Kafka は、このスループット (バイト/秒) を測定します。 30 秒を超えるスループットが構成済みの割り当て量を超えていることが検出された場合、 Kafka は、スループットが割り当て量制限内に収まる十分な遅延を計算します。

その後、 Kafka ブローカーは、標準 Kafka プロトコル応答の一部として遅延情報をクライアントに送信します。 連携クライアントは、プロトコル規約を尊重して、新しい要求を行う前にこの遅延を待機します。非連携クライアントはスロットル要求を尊重しないことがありますが、そのような場合、ブローカーはスロットルの遅延が経過するまでそのクライアントの要求を読み取りません (これにより、非連携クライアントでタイムアウトが発生する可能性があります)。

割り当て量には、以下の情報が適用されます。

  • プロデューサーとコンシューマーに対して個別の割り当て量を定義できます。
  • デフォルトでは、クライアント割り当て量は無制限です。
  • 割り当て量は、クラスターごとではなく、ブローカーごとに適用されます。
  • 割り当て量は、単一のユーザー ID を共有するすべてのクライアントに適用されます。
  • 割り当て量は「デフォルト」ユーザーに適用できるため、ユーザー固有の割り当て量が設定されていない場合でも、すべてのユーザーに適用されます。

クライアント・メトリック

Java クライアントは、以下のブローカーごとのメトリックを使用してスロットル情報も公開します。

  • produce-throttle-time-max
  • 作成-スロットル時間-平均
  • fetch-throttle-time-max
  • fetch-throttle-time-avg

クライアント割り当て量の設定

IBM® Event Streams for IBM Cloud® エンタープライズ・プランでは、 Kafka API を使用して、 Kafka V3.1.x クラスターで割り当て量を設定および記述できます。 詳しくは、 IBM® Event Streams for IBM Cloud® 管理 REST API の「 割り当て量の操作 」セクションを参照してください。

割り当て量に関するKafka 資料 を参照すると、「user」エンティティー (または「default user」) に適用されるスループット割り当て量タイプ (「producer_byte_rate」および「consumer_byte_rate」割り当て量タイプ) のみがサポートされます。

「client-id」エンティティー、「request」、および「controller-mutation」割り当て量タイプは、ユーザーが設定可能な割り当て量としてサポートされていません。 Event Streamsでは、認証済みユーザー ID は IBM Cloud® Identity and Access Management ID で表されます。 Kafka 割り当て量は Cloud Identity and Access Management ID ごとに適用されるため、API キーのグループがすべて同じ Cloud Identity and Access Management サービス ID に属している場合は、単一の割り当て量を共有できます。

Cloud Identity and Access Management (IAM) サービス ID の Cloud Identity and Access Management ID を取得するには、 IBM Cloud CLI を使用できます。

ibmcloud iam service-id ServiceId-12345678-aaaa-bbbb-cccc-1234567890ab --output json

出力は以下の例のようになります。

{

    "active":true,

    "jti":"...",

    "iam_id":"iam-ServiceId-12345678-aaaa-bbbb-cccc-1234567890ab",

    "realmId":"iam",

     ....

}

IBM Event Streams エンタープライズ・クラスターへの割り当て量のマッピング

Kafka API 割り当て量はブローカーごとですが、エンタープライズ・プランの容量は クラスターごとのスループット と呼ばれます。 したがって、ユーザーの合計を 10 MB/ 秒に制限する場合は、各ブローカーに 10/n MB/ 秒の割り当て量を適用します (n はクラスター内のブローカーの数です)。

クラスター内のブローカーの数を調べるには、 KafkaAdminClient.describeCluster 呼び出しを使用できます。

詳しくは、 Java の資料 を参照してください。

ブローカーの数は、 Apache Kafka ディストリビューションにバンドルされている kafka-configs.sh シェル・スクリプトを使用して見つけることもできます。

bin/kafka-configs.sh --command-config command-config.properties --bootstrap-server "kafka-0.blah.cloud:9093" --describe --entity-type brokers

Event Streams 許可

クライアント割り当て量を設定する権限をユーザーに付与するには、 Cloud Identity and Access Managementの「cluster」リソースに対する管理者役割を持っている必要があります。

インスタンスに対する管理者役割を持つ Event Streams UI で作成された資格情報のセットにも、クラスターに対する管理者役割があります。 したがって、割り当て量の設定に加えて、トピックの作成、削除、および変更を行うことができます。 認証済みユーザーは、クラスターに対する暗黙のリーダー役割を持ち、割り当て量を記述できます。

トピック、グループ、およびトランザクションに参加できるが、割り当て量を設定するために許可されていない資格情報のセットを作成するには、リソース・タイプ「topic」、「group」、および「txnid」に対する管理者役割と「cluster」に対するリーダー役割を持つ IAM アクセス・ポリシーをサービス ID に関連付ける必要があります。

例: the kafka-config.sh スクリプト (Apache クライアント) を使用した割り当て量の管理

  1. Kafka バイナリー ・ディストリビューション (少なくとも V3.1.0) をダウンロードします。

  2. 以下のエントリーを含むプロパティー・ファイル (以下のコマンド行の例では command-config.properties という名前) を作成します (「myapikey」を実際の API キーに置き換えます)。

sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username="token" password="myapikey";

security.protocol=SASL_SSL

sasl.mechanism=PLAIN

ssl.protocol=TLSv1.2

ssl.enabled.protocols=TLSv1.2

ssl.endpoint.identification.algorithm=HTTPS

詳しくは、Kafka API クライアントの構成を参照してください。

  1. 以下の使用例を参照してください。

    • ユーザーの割り当て量を変更します。
    bin/kafka-configs.sh --command-config command-config.properties --bootstrap-server "kafka-0.blah.cloud:9093" --alter --add-config 'producer_byte_rate=1024,
    consumer_byte_rate=2048' --entity-type users --entity-name iam-ServiceId-12345678-aaaa-bbbb-cccc-1234567890ab
    

    ユーザー iam-ServiceId-12345678-aaaa-bbbb-cccc-1234567890ab の構成の更新が完了しました。

    • ユーザーの割り当て量の記述:
    bin/kafka-configs.sh --command-config command-config.properties --bootstrap-server "kafka-0.blah.cloud:9093" --describe --entity-type users --entity-name
    iam-ServiceId-12345678-aaaa-bbbb-cccc-1234567890ab
    

    ユーザー・プリンシパル iam-ServiceId-12345678-aaaa-bbbb-cccc-1234567890ab の割り当て量構成は、 consumer_byte_rate=2048.0 および producer_byte_rate=1024.0 です。

    • ユーザーの割り当て量を削除:
    bin/kafka-configs.sh --command-config command-config.properties --bootstrap-server "kafka-0.blah.cloud:9093" --alter --delete-config       'producer_byte_rate,consumer_byte_rate' --entity-type users --entity-name iam-ServiceId-12345678-aaaa-bbbb-cccc-1234567890ab
    

    ユーザー iam-ServiceId-12345678-aaaa-bbbb-cccc-1234567890ab の構成の更新が完了しました。

    • 任意のユーザー (デフォルト・ユーザーを含む) に設定されたすべての割り当て量を記述します。
    bin/kafka-configs.sh --command-config command-config.properties --bootstrap-server "kafka-0.blah.cloud:9093" --describe --entity-type users
    
    • デフォルト・ユーザーの割り当て量を変更します。
    bin/kafka-configs.sh --command-config command-config.properties --bootstrap-server "kafka-0.blah.cloud:9093" --alter --add-config 'producer_byte_rate=1024,
    consumer_byte_rate=2048' --entity-type users --entity-default
    

    クラスター内のユーザーのデフォルト構成の更新が完了しました。

    • デフォルト・ユーザーの割り当て量を削除します。
    bin/kafka-configs.sh --command-config command-config.properties --bootstrap-server "kafka-0.blah.cloud:9093" --alter --delete-config 'producer_byte_rate,
    consumer_byte_rate' --entity-type users --entity-default
    

    クラスター内のユーザーのデフォルト構成の更新が完了しました。

例: Java API を使用したスループット割り当て量使用量の変更

以下の例は、 KafkaAdminClient.alterclientQuotas メソッドの呼び出し方法を示す短いサンプル・スニペットです。

https://kafka.apache.org/32/javadoc/org/apache/kafka/clients/admin/KafkaAdminClient.html#alterClientQuotas(java.util.Collection,org.apache.kafka.clients.admin.AlterClientQuotas
Options)

import java.util.*;

import org.apache.kafka.clients.CommonClientConfigs;

import org.apache.kafka.clients.admin.*;

import org.apache.kafka.common.config.*;

import org.apache.kafka.common.quota.*;

class Snippet {

    public static void main(String[] args) throws Exception {

        // Kafka client configuration properties.

        String mybootstrap = "...";   // from the bootstrap_endpoints of the service credentials

        String myapikey = "...";      // from the apikey of the service credentials

        Properties properties = new Properties();

        properties.put(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, mybootstrap);

        properties.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "sasl_ssl");

        properties.put(SslConfigs.SSL_ENABLED_PROTOCOLS_CONFIG, "TLSv1.2");

        properties.put(SslConfigs.SSL_PROTOCOL_CONFIG, "TLSv1.2");

        properties.put(SaslConfigs.SASL_MECHANISM, "PLAIN");

        properties.put(SaslConfigs.SASL_JAAS_CONFIG,

                "org.apache.kafka.common.security.plain.PlainLoginModule " +

                        "required username=\"token\" password=\"" + myapikey + "\";");

        AdminClient admin = AdminClient.create(properties);

        String iamID = "iam-ServiceId-12345678-aaaa-bbbb-cccc-1234567890ab"; //set iam id of target user to set quotas to

        // if null is used instead of the iamID string, the following quota alteration will be applied to the default user

        // add quotas 

        ClientQuotaEntity entity = new ClientQuotaEntity(

                Collections.singletonMap(ClientQuotaEntity.USER, iamID));

        ClientQuotaAlteration alteration = new ClientQuotaAlteration(entity,

                Arrays.asList(new ClientQuotaAlteration.Op("consumer_byte_rate", 1000.0),

                        new ClientQuotaAlteration.Op("producer_byte_rate", 1000.0)));

        admin.alterClientQuotas(Arrays.asList(alteration)).all().get();

        // describe quotas

        DescribeClientQuotasResult describeClientQuotasFuture = admin.describeClientQuotas(ClientQuotaFilter.all());

        System.out.println(describeClientQuotasFuture.entities().get());

        //remove quotas (set them to null)

        entity = new ClientQuotaEntity(Collections.singletonMap(ClientQuotaEntity.USER, iamID));

        alteration = new ClientQuotaAlteration(entity,

                Arrays.asList(new ClientQuotaAlteration.Op("consumer_byte_rate", null),

                        new ClientQuotaAlteration.Op("producer_byte_rate", null)));

        admin.alterClientQuotas(Arrays.asList(alteration)).all().get();

    }

}

IBM Cloud Activity Tracker イベント

スループット割り当て量が更新されるたびに、 Event Streams 構成イベントが生成されます。これは、 IBM Cloud® Activity Trackerでモニターできます。

Event Streamsの Activity Tracker イベントの構成について詳しくは、 Activity Tracker の資料 を参照してください。