Google Cloud Managed Service for Apache Kafka-Cluster aktualisieren

Sie können einen Google Cloud Managed Service for Apache Kafka-Cluster bearbeiten, um Attribute wie die Clustergröße (einschließlich Anzahl der vCPUs und Arbeitsspeicher), die Liste der verbundenen Subnetze, die zulässigen Quell-IP-Bereiche für öffentliche Cluster, die Konfiguration für den automatischen Neuausgleich und die mTLS-Konfiguration zu aktualisieren.

Zum Bearbeiten eines Clusters können Sie die Google Cloud Console, die Google Cloud CLI, die Clientbibliothek oder die Managed Kafka API verwenden. Sie können die Open-Source-Apache Kafka API nicht verwenden, um einen Cluster zu aktualisieren.

Für die Aktualisierung bestimmter Attribute wie die Anzahl der vCPUs und der Arbeitsspeicher muss der Dienst den Cluster möglicherweise neu starten. Der Dienst startet den Cluster jeweils einen Broker nach dem anderen neu. Während dieses Vorgangs können Anfragen an einzelne Broker fehlschlagen, aber diese Fehler sind vorübergehend. Häufig verwendete Clientbibliotheken verarbeiten diese Fehler automatisch.

Erforderliche Rollen und Berechtigungen

Bitten Sie Ihren Administrator, Ihnen die IAM-Rolle „Managed Kafka Cluster Editor (roles/managedkafka.clusterEditor)“ für Ihr Projekt zuzuweisen, um die Berechtigungen zu erhalten, die Sie zum Aktualisieren eines Clusters benötigen. Weitere Informationen zum Zuweisen von Rollen finden Sie unter Zugriff auf Projekte, Ordner und Organisationen verwalten.

Diese vordefinierte Rolle enthält die Berechtigungen, die zum Aktualisieren eines Clusters erforderlich sind. Maximieren Sie den Abschnitt Erforderliche Berechtigungen , um die notwendigen Berechtigungen anzuzeigen, die erforderlich sind:

Erforderliche Berechtigungen

Die folgenden Berechtigungen sind erforderlich, um einen Cluster zu aktualisieren:

  • Cluster bearbeiten: managedkafka.clusters.update

Sie können diese Berechtigungen auch mit benutzerdefinierten Rollen oder anderen vordefinierten Rollen erhalten.

Größe eines Clusters anpassen

Wenn Sie die Anzahl der vCPUs oder den Arbeitsspeicher eines Clusters aktualisieren, gelten die folgenden Regeln:

  • Das Verhältnis von vCPUs zu Arbeitsspeicher des Clusters muss immer zwischen 1:1 und 1:8 liegen.

  • Für jeden vorhandenen Broker müssen mindestens 1 vCPU und 1 GiB Arbeitsspeicher vorhanden sein. Die Anzahl der Broker wird nie verringert.

  • Wenn der Cluster eine benutzerdefinierte Festplattenkonfiguration, muss die Aktualisierung die Anforderungen an die Festplattenkonfiguration für den lokalen Speicher erfüllen.

  • Wenn Sie die Größe erhöhen, darf die durchschnittliche Anzahl der vCPUs und der Arbeitsspeicher pro Broker im Vergleich zu den Durchschnittswerten vor der Aktualisierung um nicht mehr als 10% sinken. Wenn Sie beispielsweise versuchen, einen Cluster von 45 vCPUs (3 Broker) auf 48 vCPUs (4 Broker) zu erhöhen, sinkt die durchschnittliche Anzahl der vCPUs pro Broker von 15 auf 12, was einer Reduzierung um 20% entspricht und damit die Grenze von 10% überschreitet.

    Wenn Sie die Anzahl der vCPUs um mehr als 10 % verringern müssen, empfehlen wir, sie in mehreren Schritten zu reduzieren. Überwachen Sie nach jeder Aktualisierung die Ressourcennutzung und gleichen Sie die Partitionen bei Bedarf neu aus.

    Wenn Sie jedoch sicher sind, dass Ihre Broker nach der Aktualisierung genügend Kapazität haben , können Sie diese Prüfung deaktivieren, indem Sie den gcloud managed-kafka clusters update Befehl mit dem allow_broker_downscale_on_cluster_upscale=true Flag ausführen. Dieses Flag signalisiert, dass Sie das potenzielle Leistungsrisiko akzeptieren.

Weitere Informationen finden Sie unter Clustergröße aktualisieren.

Konfiguration des öffentlichen Clusters

Sie können den öffentlichen Zugriff für einen vorhandenen Cluster aktivieren oder deaktivieren sowie zulässige Quell-IP-Bereiche hinzufügen oder entfernen. Weitere Informationen zu den Anforderungen und Regeln für zulässige Quell-IP-Bereiche finden Sie unter Öffentliche Cluster.

Managed Service for Apache Kafka verwendet die Cloud Next Generation Firewall, um den Zugriff auf öffentliche Cluster einzuschränken. Das Entfernen zulässiger Quell-IP-Bereiche oder das Deaktivieren des öffentlichen Zugriffs gilt nur für neue Verbindungen. Weitere Informationen finden Sie unter Auswirkungen auf vorhandenen Traffic.

Cluster bearbeiten

So bearbeiten Sie einen Cluster:

Console

  1. Rufen Sie in der Google Cloud Console die Seite Cluster auf.

Zu den Clustern

  1. Klicken Sie in der Liste der Cluster auf den Cluster, dessen Attribute Sie bearbeiten möchten.

In der Console wird die Detailseite des Clusters angezeigt.

  1. Klicken Sie auf der Cluster-Detailseite auf Bearbeiten.

  2. Bearbeiten Sie die Attribute nach Bedarf. Sie können die folgenden Attribute eines Clusters in der Console bearbeiten:

    • Arbeitsspeicher
    • vCPUs
    • Subnetz
    • Konfiguration für Neuausgleich
    • mTLS-Konfiguration
    • Labels
  3. Klicken Sie auf Speichern.

gcloud

  1. Aktivieren Sie Cloud Shell in der Google Cloud Console.

    Cloud Shell aktivieren

    Unten in der Google Cloud Console wird eine Cloud Shell Sitzung gestartet und eine Befehlszeilenaufforderung angezeigt. Cloud Shell ist eine Shell-Umgebung in der das Google Cloud CLI bereits installiert ist und Werte für Ihr aktuelles Projekt bereits festgelegt sind. Das Initialisieren der Sitzung kann einige Sekunden dauern.

  2. Ersetzen Sie folgende Werte in den Befehlsdaten:

    • PROJECT_ID: Projekt-ID.
    • LOCATION: Der Standort des Clusters.
    • CLUSTER_ID: Die ID des Clusters.
    • CPU_COUNT: Die Anzahl der vCPUs für den Cluster.
    • MEMORY: Die Größe des Arbeitsspeichers für den Cluster. Beispiel: 10GiB.
    • SUBNET_ID: Die Subnetz-ID des Subnetzes, mit dem eine Verbindung hergestellt werden soll. Beispiel: default.
    • LABELS: Die Labels, die dem Cluster zugeordnet werden sollen.
    • ALLOWED_SOURCE_IP_RANGES: Die zulässigen IPv4-CIDR-Bereiche der Quelle für den Internetzugriff auf öffentliche Cluster.

    Führen Sie den folgenden Befehl aus:

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

    Sie sollten eine Antwort ähnlich der folgenden erhalten:

    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
    
    • Verwenden Sie das Flag --no-public-cluster, um den öffentlichen Zugriff zu deaktivieren.
    • Wenn Sie das Flag --async mit Ihrem Befehl verwenden, sendet das System die Aktualisierungsanfrage und gibt sofort eine Antwort zurück, ohne auf den Abschluss des Vorgangs zu warten. Mit dem Flag --async können Sie mit anderen Aufgaben fortfahren, während die Clusteraktualisierung im Hintergrund erfolgt. Wenn Sie das Flag --async nicht verwenden, wartet das System, bis der Vorgang abgeschlossen ist, bevor es eine Antwort zurückgibt. Sie müssen warten, bis der Cluster vollständig aktualisiert wurde, bevor Sie mit anderen Aufgaben fortfahren können.

REST

Ersetzen Sie folgende Werte in den Anfragedaten:

  • PROJECT_ID: Ihre Google Cloud Projekt-ID.
  • LOCATION: Der Standort des Clusters.
  • CLUSTER_ID: Die ID des Clusters.
  • UPDATE_MASK: Die Felder, die aktualisiert werden sollen, als durch Kommas getrennte Liste vollständig qualifizierter Namen. Beispiel: capacityConfig.vcpuCount,capacityConfig.memoryBytes
  • CPU_COUNT: Die Anzahl der vCPUs für den Cluster.
  • MEMORY: Die Größe des Arbeitsspeichers für den Cluster in Byte. Beispiel: 3221225472.
  • SUBNET_ID: Die Subnetz-ID des Subnetzes, mit dem eine Verbindung hergestellt werden soll. Beispiel: default.

HTTP-Methode und URL:

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

JSON-Anfragetext:

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

Wenn Sie die Anfrage senden möchten, maximieren Sie eine der folgenden Optionen:

Sie sollten eine JSON-Antwort ähnlich wie diese erhalten:

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

Fügen Sie im Anfragetext nur die Felder ein, die Sie aktualisieren, wie im UPDATE_MASK Abfrageparameter angegeben.

  • Wenn Sie ein Subnetz hinzufügen möchten, fügen Sie networkConfigs einen neuen Eintrag im folgenden Format: projects/PROJECT_ID/regions/LOCATION/subnetworks/SUBNET_ID hinzu. Beispiel: projects/sample-project/regions/us-central1/subnetworks/default.
  • Wenn Sie den öffentlichen Zugriff aktivieren oder die zulässigen Quell-IP-Bereiche aktualisieren möchten, fügen Sie gcpConfig.accessConfig.publicClusterConfig in den UPDATE_MASK Abfrageparameter ein und geben Sie das allowedSourceIpRanges Array im Anfragetext an. Beispiel für einen Anfragetext:

    {
      "gcpConfig": {
        "accessConfig": {
          "publicClusterConfig": {
            "allowedSourceIpRanges": [
              "203.0.113.0/24"
            ]
          }
        }
      }
    }
    
  • Wenn Sie den öffentlichen Zugriff deaktivieren möchten, fügen Sie gcpConfig.accessConfig.publicClusterConfig in den UPDATE_MASK Abfrageparameter ein und übergeben Sie ein leeres JSON Objekt {} im Anfragetext (oder lassen Sie publicClusterConfig weg). Beispiel für einen Anfragetext:

    {}
    

Go

Folgen Sie der Einrichtungsanleitung für Go unter Clientbibliotheken installieren, bevor Sie dieses Beispiel anwenden. Weitere Informationen finden Sie in der Referenzdokumentation zur Managed Service for Apache Kafka Go API.

Richten Sie zur Authentifizierung bei Managed Service for Apache Kafka die Standardanmeldedaten für Anwendungen(Application Default Credentials, ADC) ein. Weitere Informationen finden Sie unter ADC für eine lokale Entwicklungsumgebung einrichten.

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

Folgen Sie der Einrichtungsanleitung für Java unter Clientbibliotheken installieren, bevor Sie dieses Beispiel anwenden. Weitere Informationen finden Sie in der Referenzdokumentation zur Managed Service for Apache Kafka Java API.

Richten Sie zur Authentifizierung bei Managed Service for Apache Kafka die Standardanmeldedaten für Anwendungen ein. Weitere Informationen finden Sie unter ADC für eine lokale Entwicklungsumgebung einrichten.


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

Folgen Sie der Einrichtungsanleitung für Python unter Clientbibliotheken installieren, bevor Sie dieses Beispiel anwenden. Weitere Informationen finden Sie in der Referenzdokumentation zur Managed Service for Apache Kafka Python API.

Richten Sie zur Authentifizierung bei Managed Service for Apache Kafka die Standardanmeldedaten für Anwendungen ein. Weitere Informationen finden Sie unter ADC für eine lokale Entwicklungsumgebung einrichten.

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

Beschränkungen

Nachdem Sie einen Managed Service for Apache Kafka-Cluster erstellt haben, können Sie die folgenden Attribute nicht mehr aktualisieren:

  • Der Clustername
  • Der Clusterstandort
  • Der Verschlüsselungstyp

Sie können den Verschlüsselungstyp zwar nicht ändern, aber Sie können Verschlüsselungsschlüssel rotieren.

Nächste Schritte

Apache Kafka® ist eine eingetragene Marke der Apache Software Foundation oder ihrer Tochtergesellschaften in den USA und/oder anderen Ländern.