Actualiza un clúster de Google Cloud Managed Service para Apache Kafka

Puedes editar un clúster de Servicio administrado para Apache Kafka de Google Cloud para actualizar propiedades como el tamaño del clúster, incluida la cantidad de CPU virtuales y la memoria, la lista de subredes conectadas, los rangos de IP de origen permitidos para clústeres públicos, la configuración de rebalanceo automático y la configuración de mTLS.

Para editar un clúster, puedes usar la Google Cloud consola de Cloud, Google Cloud CLI, la biblioteca cliente o la API de Kafka administrado. No puedes usar la API de Apache Kafka de código abierto para actualizar un clúster.

La actualización de ciertas propiedades, como la cantidad de CPU virtuales y la memoria, puede requerir que el servicio reinicie el clúster. El servicio reinicia el clúster un agente a la vez. Durante este proceso, es posible que fallen las solicitudes a agentes individuales, pero estas fallas son transitorias. Las bibliotecas cliente de uso frecuente controlan automáticamente estos errores.

Roles y permisos requeridos

Para obtener los permisos que necesitas para actualizar un clúster, pídele a tu administrador que te otorgue el rol de IAM de editor de clústeres de Kafka administrado (roles/managedkafka.clusterEditor) en tu proyecto. Para obtener más información sobre cómo otorgar roles, consulta Administra el acceso a proyectos, carpetas y organizaciones.

Este rol predefinido contiene los permisos necesarios para actualizar un clúster. Para ver los permisos exactos que son necesarios, expande la sección Permisos requeridos:

Permisos necesarios

Se requieren los siguientes permisos para actualizar un clúster:

  • Editar un clúster: managedkafka.clusters.update

También puedes obtener estos permisos con roles personalizados o otros roles predefinidos.

Cambia el tamaño de un clúster

Si actualizas la cantidad de CPU virtuales o la memoria de un clúster, se aplicarán las siguientes reglas:

  • La proporción general de CPU virtuales a memoria del clúster siempre debe permanecer entre 1:1 y 1:8.

  • Debe haber al menos 1 CPU virtual y 1 GiB de memoria para cada agente existente. La cantidad de agentes nunca disminuye.

  • Si el clúster tiene una configuración de disco personalizada, la actualización debe cumplir con los requisitos de configuración de disco para el almacenamiento local.

  • Si aumentas la escala, el promedio de CPU virtuales y memoria por agente no puede disminuir en más del 10% en comparación con los promedios anteriores a la actualización. Por ejemplo, si intentas aumentar la escala de un clúster de 45 CPU virtuales (3 agentes) a 48 CPU virtuales (4 agentes), el promedio de CPU virtuales por agente disminuye de 15 a 12, lo que representa una reducción del 20% y supera el límite del 10%.

    Si necesitas disminuir la cantidad de CPU virtuales en más del 10%, te recomendamos que lo hagas en varias etapas. Después de cada actualización, supervisa el uso de recursos y reequilibra las particiones si es necesario.

    Sin embargo, si estás seguro de que tus agentes tendrán suficiente capacidad después de la actualización, puedes inhabilitar esta verificación ejecutando el gcloud managed-kafka clusters update comando con la allow_broker_downscale_on_cluster_upscale=true marca. Esta marca indica que aceptas el riesgo potencial de rendimiento.

Para obtener más información, consulta Actualiza el tamaño del clúster.

Configuración del clúster público

Puedes habilitar o inhabilitar el acceso público para un clúster existente, así como agregar o quitar rangos de IP de origen permitidos. Para obtener más información sobre los requisitos y las reglas para los rangos de IP de origen permitidos, consulta Clústeres públicos.

Managed Service para Apache Kafka usa Cloud Next Generation Firewall para restringir el acceso a clústeres públicos. Quitar los rangos de IP de origen permitidos o inhabilitar el acceso público solo se aplica a las conexiones nuevas. Para obtener más información, consulta los efectos en el tráfico existente.

Editar un clúster

Para editar un clúster, sigue estos pasos:

Console

  1. En la Google Cloud consola de Cloud, ve a la página Clústeres.

Ir a los clústeres

  1. En la lista de clústeres, haz clic en el clúster cuyas propiedades deseas editar.

En la consola, se muestra la página de detalles del clúster.

  1. En la página de detalles del clúster, haz clic en Editar.

  2. Edita las propiedades según sea necesario. Puedes editar las siguientes propiedades de un clúster desde la consola:

    • Memoria
    • vCPUs
    • Subred
    • Configuración de rebalanceo
    • Configuración de mTLS
    • Etiquetas
  3. Haz clic en Guardar.

gcloud

  1. En la Google Cloud consola de Cloud, activa Cloud Shell.

    Activa Cloud Shell

    En la parte inferior de la Google Cloud consola de Cloud, se inicia una sesión de Cloud Shell en la que se muestra una ventana de línea de comandos. Cloud Shell es un entorno de shell con Google Cloud CLI ya instalada y con valores ya establecidos para el proyecto actual. La sesión puede tardar unos segundos en inicializarse.

  2. Antes de usar cualquiera de los datos de comando a continuación, realiza los siguientes reemplazos:

    • PROJECT_ID: Es el ID del proyecto.
    • LOCATION: Es la ubicación del clúster.
    • CLUSTER_ID: Es el ID del clúster.
    • CPU_COUNT: Es la cantidad de CPU virtuales para el clúster.
    • MEMORY: Es la cantidad de memoria para el clúster. Ejemplo: 10GiB.
    • SUBNET_ID: Es el ID de la subred a la que se conectará. Ejemplo: default.
    • LABELS: Son las etiquetas que se asociarán con el clúster.
    • ALLOWED_SOURCE_IP_RANGES: Son los rangos de CIDR de IPv4 de origen permitidos para el acceso a Internet del clúster público.

    Ejecuta el siguiente comando:

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

    Deberías recibir una respuesta similar a la que figura a continuación:

    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
    
    • Para inhabilitar el acceso público, usa la marca --no-public-cluster.
    • Si usas la marca --async con tu comando, el sistema envía la solicitud de actualización y muestra una respuesta de inmediato, sin esperar a que se complete la operación. Con la marca --async, puedes continuar con otras tareas mientras se realiza la actualización del clúster en segundo plano. Si no usas la marca --async, el sistema espera a que se complete la operación antes de mostrar una respuesta. Debes esperar hasta que el clúster esté completamente actualizado antes de continuar con otras tareas.

REST

Antes de usar cualquiera de los datos de solicitud a continuación, realiza los siguientes reemplazos:

  • PROJECT_ID: Es el ID del Google Cloud proyecto.
  • LOCATION: Es la ubicación del clúster.
  • CLUSTER_ID: Es el ID del clúster.
  • UPDATE_MASK: Son los campos que se actualizarán, como una lista separada por comas de nombres completamente calificados. Ejemplo: capacityConfig.vcpuCount,capacityConfig.memoryBytes
  • CPU_COUNT: Es la cantidad de CPU virtuales para el clúster.
  • MEMORY: Es la cantidad de memoria para el clúster, en bytes. Ejemplo: 3221225472.
  • SUBNET_ID: Es el ID de la subred a la que se conectará. Ejemplo: default.

Método HTTP y URL:

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

Cuerpo JSON de la solicitud:

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

Para enviar tu solicitud, expande una de estas opciones:

Deberías recibir una respuesta JSON similar a la que se muestra a continuación:

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

En el cuerpo de la solicitud, incluye solo los campos que estás actualizando, como se especifica en el UPDATE_MASK parámetro de consulta.

  • Para agregar una subred, agrega una entrada nueva a networkConfigs con el siguiente formato: projects/PROJECT_ID/regions/LOCATION/subnetworks/SUBNET_ID. Ejemplo: projects/sample-project/regions/us-central1/subnetworks/default.
  • Para habilitar el acceso público o actualizar los rangos de IP de origen permitidos, incluye gcpConfig.accessConfig.publicClusterConfig en el UPDATE_MASK parámetro de consulta y especifica el allowedSourceIpRanges array en el cuerpo de la solicitud. Ejemplo de cuerpo de la solicitud:

    {
      "gcpConfig": {
        "accessConfig": {
          "publicClusterConfig": {
            "allowedSourceIpRanges": [
              "203.0.113.0/24"
            ]
          }
        }
      }
    }
    
  • Para inhabilitar el acceso público, incluye gcpConfig.accessConfig.publicClusterConfig en el UPDATE_MASK parámetro de consulta y pasa un objeto JSON vacío {} en el cuerpo de la solicitud (o omite publicClusterConfig). Ejemplo de cuerpo de la solicitud:

    {}
    

Go

Antes de probar este ejemplo, sigue las instrucciones de configuración de Go en Instala las bibliotecas cliente. Para obtener más información, consulta la documentación de referencia de la API de Go de Managed Service para Apache Kafka.

Para autenticarte en Managed Service para Apache Kafka, configura las credenciales predeterminadas de la aplicación(ADC). Para obtener más información, consulta Configura ADC para un entorno de desarrollo local.

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

Antes de probar este ejemplo, sigue las instrucciones de configuración de Java en Instala las bibliotecas cliente. Para obtener más información, consulta la documentación de referencia de la API de Java de Managed Service para Apache Kafka.

Para autenticarte en Managed Service para Apache Kafka, configura las credenciales predeterminadas de la aplicación. Para obtener más información, consulta Configura ADC para un entorno de desarrollo local.


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

Antes de probar este ejemplo, sigue las instrucciones de configuración de Python en Instala las bibliotecas cliente. Para obtener más información, consulta la documentación de referencia de la API de Python de Managed Service para Apache Kafka.

Para autenticarte en Managed Service para Apache Kafka, configura las credenciales predeterminadas de la aplicación. Para obtener más información, consulta Configura ADC para un entorno de desarrollo local.

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

Limitaciones

Después de crear un clúster de Managed Service para Apache Kafka, no puedes actualizar las siguientes propiedades:

  • El nombre del clúster
  • La ubicación del clúster
  • El tipo de encriptación

Aunque no puedes cambiar el tipo de encriptación, puedes rotar las claves de encriptación.

Próximos pasos

Apache Kafka® es una marca registrada de The Apache Software Foundation o sus afiliados de Estados Unidos y otros países.