Google Cloud Managed Service for Apache Kafka-Cluster erstellen

Ein Managed Service for Apache Kafka-Cluster bietet eine Umgebung zum Speichern und Verarbeiten von Nachrichtenstreams, die in Themen organisiert sind.

Sie können einen Cluster mit der Google Cloud -Konsole, der Google Cloud CLI, der Clientbibliothek oder der Managed Kafka API erstellen. Sie können die Open-Source-Apache Kafka API nicht zum Erstellen eines Clusters verwenden.

Hinweis

Prüfen Sie, ob Sie mit Folgendem vertraut sind:

Erforderliche Rollen und Berechtigungen zum Erstellen eines Clusters

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 Erstellen 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 Erstellen eines Clusters erforderlich sind. Maximieren Sie den Abschnitt Erforderliche Berechtigungen, um die notwendigen Berechtigungen anzuzeigen:

Erforderliche Berechtigungen

Die folgenden Berechtigungen sind zum Erstellen eines Clusters erforderlich:

  • Cluster erstellen: managedkafka.clusters.create

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

Eigenschaften eines Managed Service for Apache Kafka-Clusters

Wenn Sie einen Managed Service for Apache Kafka-Cluster erstellen oder aktualisieren, müssen Sie die folgenden Eigenschaften angeben.

Clustername

Der Name oder die ID des Managed Service for Apache Kafka-Clusters, den Sie erstellen. Tipps zum Benennen von Clustern finden Sie in den Richtlinien zum Benennen einer Ressource von Managed Service for Apache Kafka. Der Name eines Clusters kann nicht geändert werden.

Standort

Der Standort, an dem Sie den Cluster erstellen. Der Standort muss eine der unterstützten Google Cloud Regionen sein. Der Standort eines Clusters kann später nicht mehr geändert werden. Eine Liste der verfügbaren Standorte finden Sie unter Managed Service for Apache Kafka-Standorte.

Kapazitätskonfiguration

Die Rechen-, Arbeitsspeicher- und Speicherkapazität des Kafka-Clusters. Beim Erstellen eines Clusters müssen Sie die folgenden Werte angeben:

  • vCPUs: Die Anzahl der vCPUs im Cluster. Mindestens 3 vCPUs pro Cluster sind erforderlich.

  • Arbeitsspeicher: Die Menge an Arbeitsspeicher, die dem Cluster zugewiesen ist. Sie müssen zwischen 1 GiB und 8 GiB pro vCPU bereitstellen.

    Wenn Sie beispielsweise einen Cluster mit 6 vCPUs erstellen, beträgt der minimale Arbeitsspeicher, den Sie dem Cluster zuweisen können, 6 GiB (1 GiB pro vCPU) und der maximale Arbeitsspeicher 48 GiB (8 GiB pro vCPU).

Der Dienst stellt automatisch Broker bereit, einschließlich nichtflüchtiger Speicher-Volumes für die Broker. Optional können Sie Benutzerdefinierter lokaler Speicher aktivieren, um die Festplattenkapazität für jeden Broker anzugeben. Weitere Informationen finden Sie unter Broker-Festplattengröße konfigurieren.

Richtlinien zum Schätzen der erforderlichen Kapazität für Ihren Cluster finden Sie unter Größe des Kafka-Clusters planen.

Netzwerkkonfiguration

Die Netzwerkkonfiguration ist eine Liste der VPC-Subnetze, in denen der Cluster zugänglich ist. Damit Clients Nachrichten senden oder empfangen können, müssen sie eines dieser Subnetze erreichen können.

Hier sind einige Richtlinien für Ihre Netzwerkkonfiguration:

  • Für einen Cluster ist mindestens ein Subnetz erforderlich. Der Höchstwert ist 10.

  • Pro Netzwerk ist für einen bestimmten Cluster genau ein Subnetz zulässig.

  • Jedes Subnetz muss sich in derselben Region wie der Cluster befinden. Das Projekt und das Netzwerk können unterschiedlich sein.

  • IP-Adressen für die Broker und den Bootstrap-Server werden in jedem Subnetz automatisch zugewiesen. Außerdem werden DNS-Einträge für diese IP-Adressen in den entsprechenden VPC-Netzwerken erstellt.

  • Wenn Sie ein Subnetz aus einem anderen Projekt hinzufügen, müssen Sie dem von Google verwalteten Dienstkonto, das dem Cluster zugeordnet ist, Berechtigungen erteilen. Weitere Informationen finden Sie unter Cluster projektübergreifend verbinden.

Nachdem Sie den Cluster erstellt haben, können Sie die Liste der Subnetze aktualisieren. Weitere Informationen zur Vernetzung finden Sie unter Vernetzung für Managed Service for Apache Kafka konfigurieren.

Öffentliche Cluster

Optional können Sie Ihren Cluster als öffentlichen Cluster konfigurieren, damit Clientanwendungen über das öffentliche Internet eine Verbindung zu Ihrem Cluster herstellen und Kafka-Vorgänge ausführen können, einschließlich des Sendens und Empfangens von Nachrichten. Wenn Sie den öffentlichen Zugriff aktivieren möchten, müssen Sie einen oder mehrere zulässige Quell-IP-Bereiche angeben.

Wenn Sie zulässige Quell-IP-Bereiche konfigurieren, gelten die folgenden Regeln:

  • Alle zulässigen Quell-IP-Bereiche müssen in IPv4-CIDR-Notation angegeben werden (z. B. 203.0.113.0/24 oder 198.51.100.5/32).

  • Die Größe jedes CIDR-Subnetzes muss zwischen /16 und /32 liegen.

  • CIDR-Bereiche dürfen sich nicht überschneiden.

  • Alle zulässigen Quell-IP-Bereiche müssen öffentlich routingfähig sein. Private IP-Bereiche wie 10.0.0.0/8 (RFC 1918) werden abgelehnt.

  • Sie können maximal 500 zulässige Quell-IP-Bereiche angeben.

  • IPv6-Adressen werden nicht unterstützt.

Bevor Sie den öffentlichen Zugriff für einen Cluster aktivieren, sollten Sie Clients mit einem öffentlichen Cluster verbinden lesen, um zu erfahren, wie die Funktion für den öffentlichen Zugriff implementiert wird.

Labels

Labels sind Schlüssel/Wert-Paare, mit denen Sie Ressourcen organisieren und identifizieren können. Mit Labels können Sie Ressourcen nach Umgebung kategorisieren. Beispiele sind "env:production" und "owner:data-engineering".

Sie können Ressourcen anhand ihrer Labels filtern und nach ihnen suchen. Angenommen, Sie haben mehrere Managed Service for Apache Kafka-Cluster für verschiedene Abteilungen. Sie können Cluster mit dem Label "department:marketing" konfigurieren und danach suchen, um den relevanten Cluster zu finden.

Konfiguration für Neuausgleich

Mit dieser Einstellung wird festgelegt, ob der Dienst Partitionen automatisch auf Broker verteilt.

Folgende Modi sind verfügbar:

  • Automatische Lastverteilung bei Hochskalierung: Wenn diese Option aktiviert ist, löst der Dienst automatisch eine Lastverteilung von Replikaten aus, wenn Sie den Cluster hochskalieren. Dieser Modus trägt zu einer gleichmäßigen Lastverteilung bei, kann aber die Leistung während des Rebalancing-Vorgangs vorübergehend beeinträchtigen.

  • Kein Neuausgleich: Wenn diese Option aktiviert ist, werden Replikate nicht automatisch neu ausgeglichen.

Weitere Informationen finden Sie unter Broker-Partitionen neu ausgleichen.

Verschlüsselung

Managed Service for Apache Kafka kann Nachrichten mitGoogle-owned and Google-managed encryption keys (Standard) oder mit vom Kunden verwalteten Verschlüsselungsschlüsseln (Customer-Managed Encryption Keys, CMEK) verschlüsseln. Jede Nachricht wird im Ruhezustand und während der Übertragung verschlüsselt. Der Verschlüsselungstyp für einen Cluster ist unveränderlich.

Standardmäßig wird Google-owned and Google-managed encryption keys verwendet. Diese Schlüssel werden vollständig von Google Cloud in seiner Infrastruktur erstellt, verwaltet und gespeichert.

CMEKs sind Verschlüsselungsschlüssel, die Sie mit Cloud Key Management Service verwalten. Mit dieser Funktion haben Sie mehr Kontrolle über die Schlüssel, die zum Verschlüsseln inaktiver Daten in unterstützten Google Cloud Diensten verwendet werden. Die Verwendung von CMEK verursacht zusätzliche Kosten im Zusammenhang mit Cloud Key Management Service. Wenn Sie CMEK verwenden, muss sich Ihr Schlüsselbund am selben Speicherort wie die Ressourcen befinden, für die Sie ihn verwenden. Weitere Informationen finden Sie unter Nachrichtenverschlüsselung konfigurieren.

mTLS-Konfiguration

Sie können mTLS optional als alternative Authentifizierungsmethode konfigurieren, bei der Clientzertifikate verwendet werden. Die Konfiguration umfasst Folgendes:

  • CA-Pools: Eine Liste mit 1 bis 10 CAS-Pools (Certificate Authority Service), denen der Cluster für die Clientauthentifizierung vertraut.

  • SSL-Prinzipalzuordnungsregeln: Eine optionale, aber empfohlene ssl.principal.mapping.rules-Broker-Eigenschaft, um lange Zertifikat-Prinzipalnamen für die Verwendung in Kafka-ACLs zu vereinfachen.

Weitere Informationen zu mTLS finden Sie unter mTLS-Authentifizierung konfigurieren.

Cluster erstellen

Bevor Sie einen Cluster erstellen, lesen Sie die Dokumentation zu Clustereigenschaften.

Das Erstellen eines Clusters dauert in der Regel 20 bis 30 Minuten.

So erstellen Sie einen Cluster:

Console

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

    Zu den Clustern

  2. Wählen Sie Erstellen aus.

    Die Seite Kafka-Cluster erstellen wird geöffnet.

  3. Geben Sie für Clustername einen String ein.

    Weitere Informationen zum Benennen eines Clusters finden Sie unter Richtlinien zum Benennen einer Managed Service for Apache Kafka-Ressource.

  4. Geben Sie unter Standort einen unterstützten Standort ein.

    Weitere Informationen zu unterstützten Standorten finden Sie unter Unterstützte Standorte für Managed Service for Apache Kafka.

  5. Geben Sie für die Kapazitätskonfiguration Werte für Arbeitsspeicher und vCPUs ein.

    Weitere Informationen zum Festlegen der Größe eines Clusters finden Sie unter Kafka-Clustergröße planen.

  6. Optional. So geben Sie die Broker-Festplattenkapazität an:

    1. Klicken Sie auf den Ein/Aus-Button Benutzerdefinierte Größe des lokalen Speichers aktivieren, um die Position „Ein“ festzulegen.

    2. Geben Sie im Feld Benutzerdefinierter lokaler Speicher die Laufwerkgröße pro Broker in GiB ein.

    Wenn Sie den Schalter Benutzerdefinierte Größe des lokalen Speichers aktivieren deaktiviert lassen, wird die Standardlaufwerkskapazität bereitgestellt. Weitere Informationen finden Sie unter Broker-Festplattengröße konfigurieren.

  7. Geben Sie unter Netzwerkkonfiguration die folgenden Details ein:

    1. Projekt: Das Projekt, in dem sich das Subnetzwerk befindet. Das Subnetz muss sich in derselben Region wie der Cluster befinden, das Projekt kann jedoch ein anderes sein.
    2. Netzwerk: Das Netzwerk, mit dem das Subnetz verbunden ist.
    3. Subnetzwerk: Der Name des Subnetzes.
    4. Subnetz-URI-Pfad: Dieses Feld wird automatisch ausgefüllt. Alternativ können Sie hier den Subnetzpfad eingeben. Der Name des Subnetzes muss das Format projects/PROJECT_ID/regions/REGION/subnetworks/SUBNET_ID haben.
    5. Klicken Sie auf Fertig.
  8. Optional: Wenn Sie weitere Subnetze hinzufügen möchten, klicken Sie auf Verbundenes Subnetz hinzufügen.

    Sie können bis zu 10 weitere Subnetze hinzufügen.

  9. Optional: Aktivieren Sie den öffentlichen Zugriff auf Ihren Cluster.

    1. Wählen Sie im Abschnitt Konfiguration des öffentlichen Clusters unter Cluster kann erreicht werden die Option In einem privaten Netzwerk und über das öffentliche Internet aus.

    2. Geben Sie im Feld Durch Kommas getrennte Liste zulässiger Quell-IPv4-Bereiche in CIDR-Notation die IP-Bereiche Ihrer Clientanwendungen in IPv4-CIDR-Notation ein (z. B. 203.0.113.0/24). IPv6-Adressen werden nicht unterstützt.

  10. Optional: Legen Sie einen kundenverwalteten Verschlüsselungsschlüssel (CMEK) fest.

    1. Wählen Sie unter Verschlüsselung die Option Cloud KMS-Schlüssel aus.

    2. Wählen Sie für Schlüsseltyp die Option Cloud KMS aus.

    3. Wählen Sie unter Vom Kunden verwalteten Schlüssel auswählen den CMEK aus oder geben Sie ihn ein.

      Wenn diese Option nicht angezeigt wird, müssen Sie die Cloud KMS API für das Projekt aktiviert haben.

    Weitere Informationen finden Sie unter Nachrichtenverschlüsselung konfigurieren.

  11. Klicken Sie auf Erstellen.

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 Eingabeaufforderung 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, bevor sie einen der Befehlsdaten verwenden:

    • 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: 3GiB.
    • SUBNET_ID: die Subnetz-ID des Subnetzes, mit dem eine Verbindung hergestellt werden soll. Beispiel: default.
    • ENCRYPTION_KEY: Die ID des kundenverwalteten Verschlüsselungsschlüssels, der für den Cluster verwendet werden soll. Dies ist optional.
    • LABELS: die Labels, die dem Cluster zugeordnet werden sollen.
    • ALLOWED_SOURCE_IP_RANGES: Die zulässigen Quell-IP-Bereiche für den öffentlichen Cluster. Dies ist optional.

    Führen Sie folgenden Befehl aus:

    Linux, macOS oder Cloud Shell

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

    Windows (PowerShell)

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

    Windows (cmd.exe)

    gcloud managed-kafka clusters create CLUSTER_ID ^
        --location=LOCATION ^
        --cpu=CPU_COUNT ^
        --memory=MEMORY ^
        --subnets=projects/PROJECT_ID/regions/LOCATION/subnetworks/SUBNET_ID ^
        --auto-rebalance ^
        --encryption-key=ENCRYPTION_KEY ^
        --async ^
        --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: create
    name: projects/PROJECT_ID/locations/LOCATION/operations/OPERATION_ID
    

    Speichere die OPERATION_ID, um progress zu tracken.

REST

Ersetzen Sie diese Werte in den folgenden Anfragedaten:

  • PROJECT_ID: Ihre Google Cloud 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 in Byte. Beispiel: 3221225472.
  • SUBNET_ID: Subnetz-ID des Subnetzes, mit dem eine Verbindung hergestellt werden soll. Beispiel: default.

HTTP-Methode und URL:

POST https://managedkafka.googleapis.com/v1/projects/PROJECT_ID/locations/LOCATION/clusters?clusterId=CLUSTER_ID

JSON-Text anfordern:

{
  "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": "create",
    "requestedCancellation": false,
    "apiVersion": "v1"
  },
  "done": false
}

Terraform

Sie können eine Terraform-Ressource verwenden, um einen Cluster zu erstellen.

resource "google_managed_kafka_cluster" "default" {
  project    = data.google_project.default.project_id # Replace this with your project ID in quotes
  cluster_id = "my-cluster-id"
  location   = "us-central1"
  capacity_config {
    vcpu_count   = 3
    memory_bytes = 3221225472
  }
  gcp_config {
    access_config {
      network_configs {
        subnet = google_compute_subnetwork.default.id
      }
    }
  }
}

Informationen zum Anwenden oder Entfernen einer Terraform-Konfiguration finden Sie unter Grundlegende Terraform-Befehle.

Go

Bevor Sie dieses Beispiel ausprobieren, folgen Sie der Einrichtungsanleitung für Go unter Clientbibliotheken installieren. 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"

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

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

	locationPath := fmt.Sprintf("projects/%s/locations/%s", projectID, region)
	clusterPath := fmt.Sprintf("%s/clusters/%s", locationPath, clusterID)

	// Memory must be between 1 GiB and 8 GiB per CPU.
	capacityConfig := &managedkafkapb.CapacityConfig{
		VcpuCount:   cpu,
		MemoryBytes: memoryBytes,
	}
	var networkConfig []*managedkafkapb.NetworkConfig
	networkConfig = append(networkConfig, &managedkafkapb.NetworkConfig{
		Subnet: subnet,
	})
	platformConfig := &managedkafkapb.Cluster_GcpConfig{
		GcpConfig: &managedkafkapb.GcpConfig{
			AccessConfig: &managedkafkapb.AccessConfig{
				NetworkConfigs: networkConfig,
			},
		},
	}
	rebalanceConfig := &managedkafkapb.RebalanceConfig{
		Mode: managedkafkapb.RebalanceConfig_AUTO_REBALANCE_ON_SCALE_UP,
	}
	cluster := &managedkafkapb.Cluster{
		Name:            clusterPath,
		CapacityConfig:  capacityConfig,
		PlatformConfig:  platformConfig,
		RebalanceConfig: rebalanceConfig,
	}

	req := &managedkafkapb.CreateClusterRequest{
		Parent:    locationPath,
		ClusterId: clusterID,
		Cluster:   cluster,
	}
	op, err := client.CreateCluster(ctx, req)
	if err != nil {
		return fmt.Errorf("client.CreateCluster got err: %w", err)
	}
	// The duration of this operation can vary considerably, typically taking 10-40 minutes.
	resp, err := op.Wait(ctx)
	if err != nil {
		return fmt.Errorf("op.Wait got err: %w", err)
	}
	fmt.Fprintf(w, "Created cluster: %s\n", resp.Name)
	return nil
}

Java

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

Richten Sie die Standardanmeldedaten für Anwendungen ein, um sich bei Managed Service for Apache Kafka zu authentifizieren. 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.RetryingFuture;
import com.google.api.gax.retrying.TimedRetryAlgorithm;
import com.google.cloud.managedkafka.v1.AccessConfig;
import com.google.cloud.managedkafka.v1.CapacityConfig;
import com.google.cloud.managedkafka.v1.Cluster;
import com.google.cloud.managedkafka.v1.CreateClusterRequest;
import com.google.cloud.managedkafka.v1.GcpConfig;
import com.google.cloud.managedkafka.v1.LocationName;
import com.google.cloud.managedkafka.v1.ManagedKafkaClient;
import com.google.cloud.managedkafka.v1.ManagedKafkaSettings;
import com.google.cloud.managedkafka.v1.NetworkConfig;
import com.google.cloud.managedkafka.v1.OperationMetadata;
import com.google.cloud.managedkafka.v1.RebalanceConfig;
import java.time.Duration;
import java.util.concurrent.ExecutionException;

public class CreateCluster {

  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";
    String subnet = "my-subnet"; // e.g. projects/my-project/regions/my-region/subnetworks/my-subnet
    int cpu = 3;
    long memoryBytes = 3221225472L; // 3 GiB
    createCluster(projectId, region, clusterId, subnet, cpu, memoryBytes);
  }

  public static void createCluster(
      String projectId, String region, String clusterId, String subnet, int cpu, long memoryBytes)
      throws Exception {
    CapacityConfig capacityConfig =
        CapacityConfig.newBuilder().setVcpuCount(cpu).setMemoryBytes(memoryBytes).build();
    NetworkConfig networkConfig = NetworkConfig.newBuilder().setSubnet(subnet).build();
    GcpConfig gcpConfig =
        GcpConfig.newBuilder()
            .setAccessConfig(AccessConfig.newBuilder().addNetworkConfigs(networkConfig).build())
            .build();
    RebalanceConfig rebalanceConfig =
        RebalanceConfig.newBuilder()
            .setMode(RebalanceConfig.Mode.AUTO_REBALANCE_ON_SCALE_UP)
            .build();
    Cluster cluster =
        Cluster.newBuilder()
            .setCapacityConfig(capacityConfig)
            .setGcpConfig(gcpConfig)
            .setRebalanceConfig(rebalanceConfig)
            .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.createClusterOperationSettings()
        .setPollingAlgorithm(timedRetryAlgorithm);

    try (ManagedKafkaClient managedKafkaClient = ManagedKafkaClient.create(
        settingsBuilder.build())) {

      CreateClusterRequest request =
          CreateClusterRequest.newBuilder()
              .setParent(LocationName.of(projectId, region).toString())
              .setClusterId(clusterId)
              .setCluster(cluster)
              .build();

      // The duration of this operation can vary considerably, typically taking between 10-40
      // minutes.
      OperationFuture<Cluster, OperationMetadata> future =
          managedKafkaClient.createClusterOperationCallable().futureCall(request);

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

      while (!future.isDone()) {
        // The pollingFuture gives us the most recent status of the operation
        RetryingFuture<OperationSnapshot> pollingFuture = future.getPollingFuture();
        OperationSnapshot currentOp = pollingFuture.getAttemptResult().get();
        System.out.printf("Polling Operation:\nName: %s\n Done: %s\n",
            currentOp.getName(),
            currentOp.isDone());
      }

      // NOTE: future.get() blocks completion until the operation is complete (isDone =  True)
      Cluster response = future.get();
      System.out.printf("Created cluster: %s\n", response.getName());
    } catch (ExecutionException e) {
      System.err.printf("managedKafkaClient.createCluster got err: %s", e.getMessage());
    }
  }
}

Python

Bevor Sie dieses Beispiel ausprobieren, folgen Sie der Anleitung zur Einrichtung von Python unter Clientbibliotheken installieren. Weitere Informationen finden Sie in der Referenzdokumentation zur Python API für Managed Service for Apache Kafka.

Richten Sie die Standardanmeldedaten für Anwendungen ein, um sich bei Managed Service for Apache Kafka zu authentifizieren. 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

# TODO(developer)
# project_id = "my-project-id"
# region = "us-central1"
# cluster_id = "my-cluster"
# subnet = "projects/my-project-id/regions/us-central1/subnetworks/default"
# cpu = 3
# memory_bytes = 3221225472

client = managedkafka_v1.ManagedKafkaClient()

cluster = managedkafka_v1.Cluster()
cluster.name = client.cluster_path(project_id, region, cluster_id)
cluster.capacity_config.vcpu_count = cpu
cluster.capacity_config.memory_bytes = memory_bytes
cluster.gcp_config.access_config.network_configs = [
    managedkafka_v1.NetworkConfig(subnet=subnet)
]
cluster.rebalance_config.mode = (
    managedkafka_v1.RebalanceConfig.Mode.AUTO_REBALANCE_ON_SCALE_UP
)

request = managedkafka_v1.CreateClusterRequest(
    parent=client.common_location_path(project_id, region),
    cluster_id=cluster_id,
    cluster=cluster,
)

try:
    operation = client.create_cluster(request=request)
    print(f"Waiting for operation {operation.operation.name} to complete...")
    # The duration of this operation can vary considerably, typically taking 10-40 minutes.
    # We can set a timeout of 3000s (50 minutes).
    response = operation.result(timeout=3000)
    print("Created cluster:", response)
except GoogleAPICallError as e:
    print(f"The operation failed with error: {e.message}")

Vorgang zur Clustererstellung überwachen

Sie können den folgenden Befehl nur ausführen, wenn Sie den Cluster mit der gcloud CLI erstellt haben.

  • Das Erstellen eines Clusters dauert in der Regel 20 bis 30 Minuten. Um den Fortschritt der Clustererstellung zu verfolgen, verwendet der Befehl gcloud managed-kafka clusters create einen Vorgang mit langer Ausführungszeit (Long-Running Operation, LRO), den Sie mit dem folgenden Befehl überwachen können:

    gcloud managed-kafka operations describe OPERATION_ID \
        --location=LOCATION
    

    Ersetzen Sie Folgendes:

    • OPERATION_ID mit dem Wert der Vorgangs-ID aus dem vorherigen Abschnitt.
    • LOCATION durch den Wert des Standorts aus dem vorherigen Abschnitt.

Fehlerbehebung

Informationen zur Fehlerbehebung beim Erstellen eines Clusters finden Sie unter Fehler beim Erstellen von Kafka-Clustern.

Nächste Schritte

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