Memperbarui cluster Google Cloud Managed Service for Apache Kafka

Anda dapat mengedit cluster Google Cloud Managed Service untuk Apache Kafka guna memperbarui properti seperti ukuran cluster, termasuk jumlah vCPU dan memori, daftar subnet yang terhubung, rentang IP sumber yang diizinkan untuk cluster publik, konfigurasi penyeimbangan ulang otomatis, dan konfigurasi mTLS.

Untuk mengedit cluster, Anda dapat menggunakan Google Cloud konsol, Google Cloud CLI, library klien, atau Managed Kafka API. Anda tidak dapat menggunakan Apache Kafka API open source untuk mengupdate cluster.

Memperbarui properti tertentu, seperti jumlah vCPU dan memori, mungkin mengharuskan layanan memulai ulang cluster. Layanan akan memulai ulang cluster satu broker dalam satu waktu. Selama proses ini, permintaan ke broker individual mungkin gagal, tetapi kegagalan ini bersifat sementara. Library klien yang umum digunakan akan otomatis menangani error ini.

Peran dan izin yang diperlukan

Untuk mendapatkan izin yang diperlukan guna mengupdate cluster, minta administrator untuk memberi Anda peran IAM Managed Kafka Cluster Editor (roles/managedkafka.clusterEditor) di project Anda. Untuk mengetahui informasi selengkapnya tentang cara memberikan peran, lihat Mengelola akses ke project, folder, dan organisasi.

Peran bawaan ini berisi izin yang diperlukan untuk mengupdate cluster. Untuk melihat izin yang benar-benar diperlukan, perluas bagian Izin yang diperlukan:

Izin yang diperlukan

Izin berikut diperlukan untuk mengupdate cluster:

  • Mengedit cluster: managedkafka.clusters.update

Anda mungkin juga bisa mendapatkan izin ini dengan peran khusus atau peran bawaan lainnya.

Mengubah ukuran cluster

Jika Anda mengupdate jumlah vCPU atau memori cluster, aturan berikut akan berlaku:

  • Rasio vCPU terhadap memori keseluruhan cluster harus selalu berada di antara 1:1 dan 1:8.

  • Harus ada minimal 1 vCPU dan 1 GiB memori untuk setiap broker yang ada. Jumlah broker tidak pernah berkurang.

  • Jika cluster memiliki konfigurasi disk kustom, update harus memenuhi persyaratan konfigurasi disk untuk penyimpanan lokal.

  • Jika Anda melakukan penskalaan, vCPU dan memori rata-rata per broker tidak boleh berkurang lebih dari 10% dibandingkan rata-rata sebelum update. Misalnya, jika Anda mencoba melakukan penskalaan cluster dari 45 vCPU (3 broker) menjadi 48 vCPU (4 broker), vCPU rata-rata per broker akan berkurang dari 15 menjadi 12, yang merupakan pengurangan sebesar 20%, sehingga melebihi batas 10%.

    Jika Anda perlu mengurangi jumlah vCPU lebih dari 10%, sebaiknya kurangi dalam beberapa tahap. Setelah setiap update, pantau penggunaan resource, dan seimbangkan kembali partisi jika diperlukan.

    Namun, jika Anda yakin bahwa broker Anda akan memiliki kapasitas yang cukup setelah update, Anda dapat menonaktifkan pemeriksaan ini dengan menjalankan gcloud managed-kafka clusters update perintah dengan flag allow_broker_downscale_on_cluster_upscale=true. Flag ini menandakan bahwa Anda menerima potensi risiko performa.

Untuk mengetahui informasi selengkapnya, lihat Mengupdate ukuran cluster.

Konfigurasi cluster publik

Anda dapat mengaktifkan atau menonaktifkan akses publik untuk cluster yang ada, serta menambahkan atau menghapus rentang IP sumber yang diizinkan. Untuk mengetahui informasi selengkapnya tentang persyaratan dan aturan untuk rentang IP sumber yang diizinkan, lihat Cluster publik.

Managed Service untuk Apache Kafka menggunakan Cloud Next Generation Firewall untuk membatasi akses ke cluster publik. Menghapus rentang IP sumber yang diizinkan atau menonaktifkan akses publik hanya berlaku untuk koneksi baru. Untuk mengetahui informasi selengkapnya, lihat efek pada traffic yang ada.

Mengedit cluster

Untuk mengedit cluster, ikuti langkah-langkah berikut:

Konsol

  1. Di Google Cloud konsol, buka halaman Cluster.

Buka Cluster

  1. Dari daftar cluster, klik cluster yang propertinya ingin Anda edit.

Konsol akan menampilkan halaman detail cluster.

  1. Di halaman detail cluster, klik Edit.

  2. Edit properti sesuai kebutuhan. Anda dapat mengedit properti cluster berikut dari konsol:

    • Memori
    • vCPUs
    • Subnet
    • Konfigurasi penyeimbangan ulang
    • Konfigurasi mTLS
    • Label
  3. Klik Simpan.

gcloud

  1. Di Google Cloud konsol, aktifkan Cloud Shell.

    Aktifkan Cloud Shell

    Di bagian bawah konsol Google Cloud , sesi Cloud Shell akan dimulai dan menampilkan prompt command line. Cloud Shell adalah lingkungan shell dengan Google Cloud CLI yang sudah terinstal, dan dengan nilai yang sudah ditetapkan untuk project Anda saat ini. Diperlukan waktu beberapa detik untuk melakukan inisialisasi pada sesi.

  2. Sebelum menggunakan salah satu data perintah di bawah, lakukan penggantian berikut:

    • PROJECT_ID: project ID.
    • LOCATION: lokasi cluster.
    • CLUSTER_ID: ID cluster.
    • CPU_COUNT: jumlah vCPU untuk cluster.
    • MEMORY: jumlah memori untuk cluster. Contoh: 10GiB.
    • SUBNET_ID: ID subnet dari subnet yang akan terhubung. Contoh: default.
    • LABELS: label yang akan dikaitkan dengan cluster.
    • ALLOWED_SOURCE_IP_RANGES: rentang CIDR IPv4 sumber yang diizinkan untuk akses internet cluster publik.

    Jalankan perintah berikut:

    Linux, macOS, atau 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

    Anda akan melihat respons seperti berikut:

    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
    
    • Untuk menonaktifkan akses publik, gunakan flag --no-public-cluster.
    • Jika Anda menggunakan flag --async dengan perintah Anda, sistem akan mengirim permintaan update dan segera menampilkan respons, tanpa menunggu operasi selesai. Dengan flag --async, Anda dapat melanjutkan tugas lain saat update cluster terjadi di latar belakang. Jika Anda tidak menggunakan flag --async, sistem akan menunggu operasi selesai sebelum menampilkan respons. Anda harus menunggu hingga cluster sepenuhnya diupdate sebelum dapat melanjutkan tugas lain.

REST

Sebelum menggunakan salah satu data permintaan, lakukan penggantian berikut:

  • PROJECT_ID: ID proyek Google Cloud Anda
  • LOCATION: lokasi cluster.
  • CLUSTER_ID: ID cluster.
  • UPDATE_MASK: kolom yang akan diupdate, sebagai daftar yang dipisahkan koma dari nama yang sepenuhnya memenuhi syarat. Contoh: capacityConfig.vcpuCount,capacityConfig.memoryBytes
  • CPU_COUNT: jumlah vCPU untuk cluster.
  • MEMORY: jumlah memori untuk cluster, dalam byte. Contoh: 3221225472.
  • SUBNET_ID: ID subnet dari subnet yang akan terhubung. Contoh: default.

Metode HTTP dan URL:

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

Meminta isi JSON:

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

Untuk mengirim permintaan Anda, perluas salah satu opsi berikut:

Anda akan melihat respons JSON seperti berikut:

{
  "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
}

Dalam isi permintaan, sertakan hanya kolom yang Anda update, seperti yang ditentukan dalam parameter kueri UPDATE_MASK.

  • Untuk menambahkan subnet, tambahkan entri baru ke networkConfigs dalam format berikut: projects/PROJECT_ID/regions/LOCATION/subnetworks/SUBNET_ID. Contoh: projects/sample-project/regions/us-central1/subnetworks/default.
  • Untuk mengaktifkan akses publik atau mengupdate rentang IP sumber yang diizinkan, sertakan gcpConfig.accessConfig.publicClusterConfig dalam parameter kueri UPDATE_MASK dan tentukan array allowedSourceIpRanges dalam isi permintaan. Contoh isi permintaan:

    {
      "gcpConfig": {
        "accessConfig": {
          "publicClusterConfig": {
            "allowedSourceIpRanges": [
              "203.0.113.0/24"
            ]
          }
        }
      }
    }
    
  • Untuk menonaktifkan akses publik, sertakan gcpConfig.accessConfig.publicClusterConfig dalam parameter kueri UPDATE_MASK dan teruskan objek JSON kosong {} dalam isi permintaan (atau hapuspublicClusterConfig). Contoh isi permintaan:

    {}
    

Go

Sebelum mencoba contoh ini, ikuti petunjuk penyiapan Go di Menginstal library klien. Untuk mengetahui informasi selengkapnya, lihat Dokumentasi referensi API Managed Service untuk Apache Kafka Go.

Untuk melakukan autentikasi ke Managed Service untuk Apache Kafka, siapkan Kredensial Default Aplikasi(ADC). Untuk mengetahui informasi selengkapnya, lihat Menyiapkan ADC untuk lingkungan pengembangan lokal.

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

Sebelum mencoba contoh ini, ikuti petunjuk penyiapan Java di Menginstal library klien. Untuk mengetahui informasi selengkapnya, lihat Dokumentasi referensi API Managed Service untuk Apache Kafka Java.

Untuk melakukan autentikasi ke Managed Service untuk Apache Kafka, siapkan Kredensial Default Aplikasi. Untuk mengetahui informasi selengkapnya, lihat Menyiapkan ADC untuk lingkungan pengembangan lokal.


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

Sebelum mencoba contoh ini, ikuti petunjuk penyiapan Python di Menginstal library klien. Untuk mengetahui informasi selengkapnya, lihat Dokumentasi referensi Managed Service untuk Apache Kafka Python API.

Untuk melakukan autentikasi ke Managed Service untuk Apache Kafka, siapkan Kredensial Default Aplikasi. Untuk mengetahui informasi selengkapnya, lihat Menyiapkan ADC untuk lingkungan pengembangan lokal.

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}")

Batasan

Setelah membuat cluster Managed Service untuk Apache Kafka, Anda tidak dapat mengupdate properti berikut:

  • Nama cluster
  • Lokasi cluster
  • Jenis enkripsi

Meskipun Anda tidak dapat mengubah jenis enkripsi, Anda dapat merotasi kunci enkripsi.

Apa langkah selanjutnya?

Apache Kafka® adalah merek dagang terdaftar dari The Apache Software Foundation atau afiliasinya di Amerika Serikat dan/atau negara lain.