Google Cloud Managed Service for Apache Kafka クラスタを更新する

Google Cloud Managed Service for Apache Kafka クラスタを編集して、クラスタサイズ(vCPU 数とメモリを含む)、接続されているサブネットのリスト、パブリック クラスタで許可される送信元 IP 範囲、自動再分散の構成、mTLS の構成などのプロパティを更新できます。

クラスタを編集するには、 Google Cloud コンソール、Google Cloud CLI、 クライアント ライブラリ、Managed Kafka API を使用します。オープンソースの Apache Kafka API を使用してクラスタを更新することはできません。

vCPU 数やメモリなどの特定のプロパティを更新すると、サービスでクラスタの再起動が必要になることがあります。サービスは、ブローカーごとに 1 つずつクラスタを再起動します。このプロセス中に、個々のブローカーへのリクエストが失敗することがありますが、これらの障害は一時的なものです。一般的に使用されるクライアント ライブラリは、これらのエラーを自動的に処理します。

必要なロールと権限

クラスタを更新するために必要な権限を取得するには、プロジェクトに対するManaged Kafka クラスタ編集者 roles/managedkafka.clusterEditor)IAM ロールを付与するよう管理者に依頼してください。ロールの付与については、プロジェクト、フォルダ、組織に対するアクセス権の管理をご覧ください。

この事前定義ロールには クラスタの更新に必要な権限が含まれています。必要とされる正確な権限については、「必要な権限」セクションを開いてご確認ください。

必要な権限

クラスタを更新するには、次の権限が必要です。

  • クラスタの編集: managedkafka.clusters.update

カスタムロールや他の事前定義ロールを使用して、これらの権限を取得することもできます。

クラスタのサイズを変更する

クラスタの vCPU 数またはメモリを更新する場合は、次のルールが適用されます。

  • クラスタの全体的な vCPU 対メモリの比率は、常に 1:1 ~ 1:8 の範囲内にする必要があります。

  • 既存のブローカーごとに、少なくとも 1 つの vCPU と 1 GiB のメモリが必要です。ブローカーの数が減ることはありません。

  • クラスタに カスタム ディスク構成がある場合、 更新はローカル ストレージの ディスク構成要件 を満たす必要があります。

  • スケールアップする場合、ブローカーあたりの平均 vCPU とメモリは、更新前の平均値と比較して 10% 以上減少させることはできません。たとえば、クラスタを 45 個の vCPU(3 個のブローカー)から 48 個の vCPU(4 個のブローカー)にスケールアップしようとすると、ブローカーあたりの平均 vCPU が 15 から 12 に減少し、20% の削減となり、10% の上限を超えます。

    vCPU 数を 10% 以上減らす必要がある場合は、段階的に減らすことをおすすめします。更新ごとにリソース使用率をモニタリングし、必要に応じてパーティションを再分散します。

    ただし、更新後にブローカーに十分な容量があることが確実な場合は、 コマンドを gcloud managed-kafka clusters update allow_broker_downscale_on_cluster_upscale=trueフラグを指定して実行することで、このチェックを無効にできます。このフラグは、潜在的なパフォーマンス リスクを受け入れることを示します。

詳細については、クラスタの サイズを更新するをご覧ください。

パブリック クラスタの構成

既存のクラスタのパブリック アクセスを有効または無効にしたり、許可された送信元 IP 範囲を追加または削除したりできます。許可された送信元 IP 範囲の要件 とルールについては、パブリック クラスタをご覧ください。

Managed Service for Apache Kafka は、 Cloud Next Generation Firewall を使用してパブリック クラスタへのアクセスを制限します。許可された送信元 IP 範囲を削除したり、パブリック アクセスを無効にしたりすると、新しい接続にのみ適用されます。詳細については、既存のトラフィックへの 影響をご覧ください。

クラスタの編集

クラスタを編集する手順は次のとおりです。

コンソール

  1. Google Cloud コンソールで、[クラスタ] ページに移動します。

[クラスタ] に移動

  1. クラスタのリストから、プロパティを編集するクラスタをクリックします。

コンソールにクラスタの詳細ページが表示されます。

  1. クラスタの詳細ページで、[編集] をクリックします。

  2. 必要に応じてプロパティを編集します。コンソールからクラスタの次のプロパティを編集できます。

    • メモリ
    • vCPUs
    • サブネット
    • 再分散の構成
    • mTLS の構成
    • ラベル
  3. [保存] をクリックします。

gcloud

  1. コンソール Google Cloud で Cloud Shell をアクティブにします。

    Cloud Shell をアクティブにする

    コンソールの下部にある Google Cloud Cloud Shell セッションが開始し、コマンドライン プロンプトが表示されます。Cloud Shell はシェル環境です 。Google Cloud CLI がすでにインストールされており、現在のプロジェクトの値もすでに設定されています 。セッションが初期化されるまで数秒かかることがあります。

  2. 後述のコマンドデータを使用する前に、 次のように置き換えます。

    • PROJECT_ID: プロジェクト ID。
    • LOCATION: クラスタのロケーション。
    • CLUSTER_ID: クラスタの ID。
    • CPU_COUNT: クラスタの vCPU 数。
    • MEMORY: クラスタのメモリ量。例: 10GiB
    • SUBNET_ID: 接続するサブネットのサブネット ID。例: default
    • LABELS: クラスタに関連付けるラベル。
    • ALLOWED_SOURCE_IP_RANGES: パブリック クラスタのインターネット アクセスで許可される送信元 IPv4 CIDR 範囲。

    次のコマンドを実行します。

    Linux、macOS、Cloud Shell

    gcloud managed-kafka clusters update CLUSTER_ID \
        --location=LOCATION \
        --cpu=CPU_COUNT \
        --memory=MEMORY \
        --subnets=projects/PROJECT_ID/regions/LOCATION/subnetworks/SUBNET_ID \
        --auto-rebalance \
        --labels=LABELS \
        --public-cluster \
        --allowed-source-ip-ranges=ALLOWED_SOURCE_IP_RANGES

    Windows(PowerShell)

    gcloud managed-kafka clusters update CLUSTER_ID `
        --location=LOCATION `
        --cpu=CPU_COUNT `
        --memory=MEMORY `
        --subnets=projects/PROJECT_ID/regions/LOCATION/subnetworks/SUBNET_ID `
        --auto-rebalance `
        --labels=LABELS `
        --public-cluster `
        --allowed-source-ip-ranges=ALLOWED_SOURCE_IP_RANGES

    Windows(cmd.exe)

    gcloud managed-kafka clusters update CLUSTER_ID ^
        --location=LOCATION ^
        --cpu=CPU_COUNT ^
        --memory=MEMORY ^
        --subnets=projects/PROJECT_ID/regions/LOCATION/subnetworks/SUBNET_ID ^
        --auto-rebalance ^
        --labels=LABELS ^
        --public-cluster ^
        --allowed-source-ip-ranges=ALLOWED_SOURCE_IP_RANGES

    次のようなレスポンスが返されます。

    done: false
    metadata:
      '@type': type.googleapis.com/google.cloud.managedkafka.v1.OperationMetadata
      apiVersion: v1
      createTime: 'CREATE_TIME'
      requestedCancellation: false
      target: projects/PROJECT_ID/locations/LOCATION/clusters/CLUSTER_ID
      verb: update
    name: projects/PROJECT_ID/locations/LOCATION/operations/OPERATION_ID
    
    • パブリック アクセスを無効にするには、--no-public-cluster フラグを使用します。
    • コマンドで --async フラグを使用すると、オペレーションの完了を待たずに更新リクエストが送信され、すぐにレスポンスが返されます。--async フラグを使用すると、クラスタの更新がバックグラウンドで行われている間に、他のタスクを続行できます。--async フラグを使用しない場合、レスポンスが返される前にオペレーションが完了するまで待機します。他のタスクを続行する前に、クラスタが完全に更新されるまで待つ必要があります。

REST

リクエスト データを使用する前に、 次のように置き換えます:

  • PROJECT_ID: あなたの Google Cloud プロジェクト ID
  • LOCATION: クラスタのロケーション。
  • CLUSTER_ID: クラスタの ID。
  • UPDATE_MASK: 更新するフィールド( 完全修飾名のカンマ区切りリスト)。例: capacityConfig.vcpuCount,capacityConfig.memoryBytes
  • CPU_COUNT: クラスタの vCPU 数。
  • MEMORY: クラスタのメモリ量(バイト単位)。例: 3221225472
  • SUBNET_ID: 接続するサブネットのサブネット ID。例: default

HTTP メソッドと URL:

PATCH https://managedkafka.googleapis.com/v1/projects/PROJECT_ID/locations/LOCATION/clusters/CLUSTER_ID?updateMask=UPDATE_MASK

リクエストの JSON 本文:

{
  "capacityConfig": {
    "vcpuCount": CPU_COUNT,
    "memoryBytes": MEMORY
  },
  "gcpConfig": {
    "accessConfig": {
      "networkConfigs": [
        {
          "subnet": "projects/PROJECT_ID/regions/LOCATION/subnetworks/SUBNET_ID"
        }
      ]
    }
  }
}

リクエストを送信するには、次のいずれかのオプションを展開します。

次のような JSON レスポンスが返されます。

{
  "name": "projects/PROJECT_ID/locations/LOCATION/operations/OPERATION_ID",
  "metadata": {
    "@type": "type.googleapis.com/google.cloud.managedkafka.v1.OperationMetadata",
    "createTime": "CREATE_TIME",
    "target": "projects/PROJECT_ID/locations/LOCATION/clusters/CLUSTER_ID",
    "verb": "update",
    "requestedCancellation": false,
    "apiVersion": "v1"
  },
  "done": false
}

リクエストの本文には、 クエリ パラメータUPDATE_MASKで指定した、更新するフィールドのみを含めます。

  • サブネットを追加するには、次の 形式で新しいエントリを networkConfigs に追加します。 projects/PROJECT_ID/regions/LOCATION/subnetworks/SUBNET_ID. 例: projects/sample-project/regions/us-central1/subnetworks/default
  • パブリック アクセスを有効にするか、許可された送信元 IP 範囲を更新するには、 gcpConfig.accessConfig.publicClusterConfig クエリ パラメータに UPDATE_MASK を含め、リクエストの本文で allowedSourceIpRanges 配列を指定します。 リクエスト本文の例:

    {
      "gcpConfig": {
        "accessConfig": {
          "publicClusterConfig": {
            "allowedSourceIpRanges": [
              "203.0.113.0/24"
            ]
          }
        }
      }
    }
    
  • パブリック アクセスを無効にするには、 gcpConfig.accessConfig.publicClusterConfigUPDATE_MASK クエリ パラメータに含め、リクエストの本文で空の JSON オブジェクト {} を渡します(または publicClusterConfig を省略します)。 リクエスト本文の例:

    {}
    

Go

このサンプルを試す前に、 クライアント ライブラリをインストールするにある Go の設定手順を行ってください。詳細については、 Managed Service for Apache Kafka Go API リファレンス ドキュメントをご覧ください

Managed Service for Apache Kafka で認証を行うには、アプリケーションのデフォルト認証情報(ADC)を設定します。 詳細については、 ローカル開発環境の ADC の設定をご覧ください。

import (
	"context"
	"fmt"
	"io"

	"cloud.google.com/go/managedkafka/apiv1/managedkafkapb"
	"google.golang.org/api/option"
	"google.golang.org/protobuf/types/known/fieldmaskpb"

	managedkafka "cloud.google.com/go/managedkafka/apiv1"
)

func updateCluster(w io.Writer, projectID, region, clusterID string, memory int64, opts ...option.ClientOption) error {
	// projectID := "my-project-id"
	// region := "us-central1"
	// clusterID := "my-cluster"
	// memoryBytes := 4221225472
	ctx := context.Background()
	client, err := managedkafka.NewClient(ctx, opts...)
	if err != nil {
		return fmt.Errorf("managedkafka.NewClient got err: %w", err)
	}
	defer client.Close()

	clusterPath := fmt.Sprintf("projects/%s/locations/%s/clusters/%s", projectID, region, clusterID)
	capacityConfig := &managedkafkapb.CapacityConfig{
		MemoryBytes: memory,
	}
	cluster := &managedkafkapb.Cluster{
		Name:           clusterPath,
		CapacityConfig: capacityConfig,
	}
	paths := []string{"capacity_config.memory_bytes"}
	updateMask := &fieldmaskpb.FieldMask{
		Paths: paths,
	}

	req := &managedkafkapb.UpdateClusterRequest{
		UpdateMask: updateMask,
		Cluster:    cluster,
	}
	op, err := client.UpdateCluster(ctx, req)
	if err != nil {
		return fmt.Errorf("client.UpdateCluster got err: %w", err)
	}
	resp, err := op.Wait(ctx)
	if err != nil {
		return fmt.Errorf("op.Wait got err: %w", err)
	}
	fmt.Fprintf(w, "Updated cluster: %#v\n", resp)
	return nil
}

Java

このサンプルを試す前に、 クライアント ライブラリをインストールするにある Java の設定手順を行ってください。詳細については、 Managed Service for Apache Kafka Java API リファレンス ドキュメントをご覧ください。

Managed Service for Apache Kafka で認証を行うには、アプリケーションのデフォルト認証情報を設定します。 詳細については、 ローカル開発環境の ADC の設定をご覧ください。


import com.google.api.gax.longrunning.OperationFuture;
import com.google.api.gax.longrunning.OperationSnapshot;
import com.google.api.gax.longrunning.OperationTimedPollAlgorithm;
import com.google.api.gax.retrying.RetrySettings;
import com.google.api.gax.retrying.TimedRetryAlgorithm;
import com.google.cloud.managedkafka.v1.CapacityConfig;
import com.google.cloud.managedkafka.v1.Cluster;
import com.google.cloud.managedkafka.v1.ClusterName;
import com.google.cloud.managedkafka.v1.ManagedKafkaClient;
import com.google.cloud.managedkafka.v1.ManagedKafkaSettings;
import com.google.cloud.managedkafka.v1.OperationMetadata;
import com.google.cloud.managedkafka.v1.UpdateClusterRequest;
import com.google.protobuf.FieldMask;
import java.time.Duration;
import java.util.concurrent.ExecutionException;

public class UpdateCluster {

  public static void main(String[] args) throws Exception {
    // TODO(developer): Replace these variables before running the example.
    String projectId = "my-project-id";
    String region = "my-region"; // e.g. us-east1
    String clusterId = "my-cluster";
    long memoryBytes = 25769803776L; // 24 GiB
    updateCluster(projectId, region, clusterId, memoryBytes);
  }

  public static void updateCluster(
      String projectId, String region, String clusterId, long memoryBytes) throws Exception {
    CapacityConfig capacityConfig = CapacityConfig.newBuilder().setMemoryBytes(memoryBytes).build();
    Cluster cluster =
        Cluster.newBuilder()
            .setName(ClusterName.of(projectId, region, clusterId).toString())
            .setCapacityConfig(capacityConfig)
            .build();
    FieldMask updateMask = FieldMask.newBuilder().addPaths("capacity_config.memory_bytes").build();

    // Create the settings to configure the timeout for polling operations
    ManagedKafkaSettings.Builder settingsBuilder = ManagedKafkaSettings.newBuilder();
    TimedRetryAlgorithm timedRetryAlgorithm = OperationTimedPollAlgorithm.create(
        RetrySettings.newBuilder()
            .setTotalTimeoutDuration(Duration.ofHours(1L))
            .build());
    settingsBuilder.updateClusterOperationSettings()
        .setPollingAlgorithm(timedRetryAlgorithm);

    try (ManagedKafkaClient managedKafkaClient = ManagedKafkaClient.create(
        settingsBuilder.build())) {
      UpdateClusterRequest request =
          UpdateClusterRequest.newBuilder().setUpdateMask(updateMask).setCluster(cluster).build();
      OperationFuture<Cluster, OperationMetadata> future =
          managedKafkaClient.updateClusterOperationCallable().futureCall(request);

      // Get the initial LRO and print details. CreateCluster contains sample code for polling logs.
      OperationSnapshot operation = future.getInitialFuture().get();
      System.out.printf("Cluster update started. Operation name: %s\nDone: %s\nMetadata: %s\n",
          operation.getName(),
          operation.isDone(),
          future.getMetadata().get().toString());

      Cluster response = future.get();
      System.out.printf("Updated cluster: %s\n", response.getName());
    } catch (ExecutionException e) {
      System.err.printf("managedKafkaClient.updateCluster got err: %s", e.getMessage());
    }
  }
}

Python

このサンプルを試す前に、 クライアント ライブラリをインストールするにある Python の設定手順を行ってください。詳細については、 Managed Service for Apache Kafka Python API リファレンス ドキュメントをご覧ください。

Managed Service for Apache Kafka で認証を行うには、アプリケーションのデフォルト認証情報を設定します。 詳細については、 ローカル開発環境の ADC の設定をご覧ください。

from google.api_core.exceptions import GoogleAPICallError
from google.cloud import managedkafka_v1
from google.protobuf import field_mask_pb2

# TODO(developer)
# project_id = "my-project-id"
# region = "us-central1"
# cluster_id = "my-cluster"
# memory_bytes = 4295000000

client = managedkafka_v1.ManagedKafkaClient()

cluster = managedkafka_v1.Cluster()
cluster.name = client.cluster_path(project_id, region, cluster_id)
cluster.capacity_config.memory_bytes = memory_bytes
update_mask = field_mask_pb2.FieldMask()
update_mask.paths.append("capacity_config.memory_bytes")

# For a list of editable fields, one can check https://cloud.google.com/managed-kafka/docs/create-cluster#properties.
request = managedkafka_v1.UpdateClusterRequest(
    update_mask=update_mask,
    cluster=cluster,
)

try:
    operation = client.update_cluster(request=request)
    print(f"Waiting for operation {operation.operation.name} to complete...")
    response = operation.result()
    print("Updated cluster:", response)
except GoogleAPICallError as e:
    print(f"The operation failed with error: {e.message}")

制限事項

Managed Service for Apache Kafka クラスタを作成した後は、次のプロパティを更新できません。

  • クラスタ名
  • クラスタのロケーション
  • 暗号化のタイプ

暗号化のタイプは変更できませんが、暗号鍵 をローテーションすることはできます。

次のステップ

Apache Kafka® は、Apache Software Foundation または米国その他の諸国における関連会社の商標です。