更新 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 數量和記憶體) 時,服務可能需要重新啟動叢集。這項服務會一次重新啟動叢集中的一個代理程式。在此期間,對個別代理商的要求可能會失敗,但這些失敗是暫時性的。常用的用戶端程式庫會自動處理這些錯誤。

必要角色和權限

如要取得更新叢集所需的權限,請要求管理員授予您專案的代管 Kafka 叢集編輯者 (roles/managedkafka.clusterEditor) IAM 角色。如要進一步瞭解如何授予角色,請參閱「管理專案、資料夾和組織的存取權」。

這個預先定義的角色具備更新叢集所需的權限。如要查看確切的必要權限,請展開「Required permissions」(必要權限) 部分:

所需權限

如要更新叢集,您必須具備下列權限:

  • 編輯叢集: 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% 的限制。

    如果需要減少超過 10% 的 vCPU 數量,建議分階段減少。每次更新後,請監控資源用量,並視需要重新平衡分割區。

    不過,如果您確信更新後代理程式仍有足夠容量,可以執行 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 控制台的「Clusters」(叢集) 頁面。

前往「Clusters」(叢集)

  1. 在叢集清單中,按一下要編輯屬性的叢集。

主控台會顯示叢集詳細資料頁面。

  1. 在叢集詳細資料頁面中,按一下「編輯」

  2. 視需要編輯屬性。您可以透過控制台編輯叢集的下列屬性:

    • 記憶體
    • vCPU
    • 子網路
    • 重新平衡設定
    • mTLS 設定
    • 標籤
  3. 按一下 [儲存]

gcloud

  1. 在 Google Cloud 控制台中啟用 Cloud Shell。

    啟用 Cloud Shell

    控制台底部會開啟 Cloud Shell 工作階段,並顯示指令列提示。 Google Cloud 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 方法和網址:

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 範圍,請在 UPDATE_MASK 查詢參數中加入 gcpConfig.accessConfig.publicClusterConfig,並在要求內文中指定 allowedSourceIpRanges 陣列。要求主體範例:

    {
      "gcpConfig": {
        "accessConfig": {
          "publicClusterConfig": {
            "allowedSourceIpRanges": [
              "203.0.113.0/24"
            ]
          }
        }
      }
    }
    
  • 如要停用公開存取權,請在 UPDATE_MASK 查詢參數中加入 gcpConfig.accessConfig.publicClusterConfig,並在要求內容中傳遞空白 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® 是 The Apache Software Foundation 或其關聯企業在美國與/或其他國家/地區的註冊商標。