Crea un clúster de Connect

Un clúster de Connect proporciona un entorno para los conectores que ayuda a transferir datos de implementaciones existentes de Kafka a un clúster de Google Cloud Managed Service para Apache Kafka o a transferir datos de un clúster de Managed Service para Apache Kafka a otro servicio Google Cloud o a otro clúster de Kafka. El clúster de Kafka secundario puede ser otro clúster de Google Cloud Managed Service para Apache Kafka, uno autoadministrado o uno local.

Antes de comenzar

Asegúrate de haber creado ya un clúster de Managed Service para Apache Kafka. Necesitas el nombre del clúster de Managed Service para Apache Kafka al que se adjuntará el clúster de Connect.

Cada clúster de Connect está asociado con un clúster de Managed Service para Apache Kafka. Este clúster almacena el estado de los conectores que se ejecutan en el clúster de Connect.

Roles y permisos obligatorios para crear un clúster de Connect

Para obtener los permisos que necesitas para crear un clúster de Connect, pídele a tu administrador que te otorgue el rol de IAM Editor de clústeres de Kafka Connect administrados (roles/managedkafka.connectClusterEditor) 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 crear un clúster de Connect. Para ver los permisos exactos que son necesarios, expande la sección Permisos requeridos:

Permisos necesarios

Se requieren los siguientes permisos para crear un clúster de Connect:

  • Otorga el permiso para crear un clúster de Connect en la ubicación especificada: managedkafka.connectClusters.create

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

Principales de LCA requeridos

De forma predeterminada, los clústeres de Managed Service para Apache Kafka permiten que el clúster de Connect acceda a los recursos si no se configuran ACL. Para ello, se debe configurar allow.everyone.if.no.acl.found como true, que es el parámetro de configuración predeterminado.

Sin embargo, si el clúster de Managed Service para Apache Kafka tiene ACL configuradas, el clúster de Connect no tendrá automáticamente los permisos de lectura y escritura para los recursos. Debes otorgarlos de forma manual.

La cuenta de servicio del clúster de Connect que se usa como principal en las ACL sigue este formato: User:service-{consumer project number}@gcp-sa-managedkafka.iam.gserviceaccount.com.

Si configuraste ACL en tu clúster de Kafka, otorga permisos de lectura y escritura a los temas y permisos de lectura a los grupos de consumidores del clúster de Connect con los siguientes comandos:

/bin/kafka-acls.sh \
    --bootstrap-server BOOTSTRAP_ADDR \
    --command-config PATH_TO_CLIENT_PROPERTIES \
    --add \
    --allow-principal User:service-{consumer project number}@gcp-sa-managedkafka.iam.gserviceaccount.com \
    --operation READ --operation WRITE --topic *
/bin/kafka-acls.sh \
    --bootstrap-server BOOTSTRAP_ADDR \
    --command-config PATH_TO_CLIENT_PROPERTIES \
    --add \
    --allow-principal User:service-{consumer project number}@gcp-sa-managedkafka.iam.gserviceaccount.com \
    --operation READ --group *

Para obtener más información sobre estos comandos, consulta Configura LCA de Apache Kafka para un control de acceso detallado.

Crea un clúster de Connect en otro proyecto

Managed Service para Apache Kafka usa un agente de servicio para acceder a los recursos deGoogle Cloud . El agente de servicio está asociado con el proyecto en el que creas el clúster.

Si creas un clúster de Connect en un proyecto diferente del clúster de Managed Service para Apache Kafka, el clúster de Connect y el clúster de Kafka usarán los agentes de servicio asociados con sus respectivos proyectos. En ese caso, el agente de servicio del clúster de Connect necesita permiso para acceder a los recursos Google Cloud en el proyecto del clúster de Kafka.

Para otorgar los permisos necesarios, asigna el rol de Agente de servicio de Kafka administrado al agente de servicio del clúster de Connect en el proyecto del clúster de Kafka. Por ejemplo, si creas un clúster de Kafka en el proyecto kafka-project y un clúster de Connect en el proyecto connect-project, otorga el rol de agente de servicio de Kafka administrado en kafka-project al agente de servicio asociado con connect-project.

La dirección de correo electrónico del agente de servicio tiene el siguiente formato: service-PROJECT_NUMBER@gcp-sa-managedkafka.iam.gserviceaccount.com, en el que PROJECT_NUMBER es el número del proyecto. Para obtener más información sobre cómo otorgar el rol, consulta Crea y otorga roles a los agentes de servicio.

Propiedades de un clúster de Connect

En esta sección, se describen las propiedades de un clúster de Connect.

Nombre del clúster de Connect

Es el nombre del clúster de Connect que estás creando. Si necesitas ayuda para asignarle un nombre a un clúster de Connect, consulta los Lineamientos para asignarles nombres a los recursos de Managed Service para Apache Kafka. Los nombres de clústeres son inmutables.

Clúster principal de Kafka

Es el clúster de Managed Service para Apache Kafka asociado a tu clúster de Connect. Este clúster asociado (clúster principal) almacena el estado de los conectores que se ejecutan en el clúster de Connect. Por lo general, el clúster principal de Managed Service para Apache Kafka también funciona como destino de todos los conectores de origen y como entrada de todos los conectores receptores que se ejecutan en el clúster de Connect.

Un solo clúster de Managed Service para Apache Kafka puede tener varios clústeres de Connect. Si eliges un clúster de Managed Service para Apache Kafka en un proyecto diferente, asegúrate de que se hayan configurado los permisos adecuados.

No puedes actualizar a un clúster de Kafka diferente después de crear el clúster de Connect.

Beneficios de la colocación en la misma región para la latencia y los costos de red

La colocación conjunta de tus clústeres de Managed Service para Apache Kafka y Connect en la misma región reduce la latencia y los costos de red. Por ejemplo, supongamos que tu clúster de Managed Service para Apache Kafka está en region-a y que usas un conector receptor para escribir datos desde este clúster de Managed Service para Apache Kafka (fuente) en una tabla de BigQuery (receptor) que también está en region-a. Si implementas tu clúster de Connect en region-a, esta elección de implementación minimiza la latencia de la operación de escritura de BigQuery y elimina los costos de transferencia de red entre regiones entre el clúster de Managed Service para Apache Kafka y el clúster de Connect.

Consideraciones de latencia y costos en varios sistemas

Kafka Connect usa conectores para transferir datos entre sistemas. Un lado del conector siempre interactúa con un clúster de Managed Service para Apache Kafka. Un solo clúster de Kafka Connect puede ejecutar varios conectores, cada uno de los cuales actúa como una fuente (extrae datos de un sistema) o un receptor (envía datos a un sistema).

Si bien un clúster de Connect en la misma región que el clúster de Managed Service para Apache Kafka se beneficia de una menor latencia de comunicación entre ellos, cada conector también interactúa con otro sistema, como una tabla de BigQuery o otro clúster de Kafka. Incluso si el clúster de Connect y el clúster de Managed Service para Apache Kafka se encuentran en la misma ubicación, ese otro sistema podría estar en una región diferente. Esto genera una mayor latencia y un mayor costo. La latencia general de la canalización depende de las ubicaciones de los tres sistemas: el clúster de Managed Service para Apache Kafka, el clúster de Connect y el sistema de origen o destino.

Por ejemplo, si tu clúster de Managed Service para Apache Kafka está en region-a, tu clúster de Connect en region-b y usas un conector de Cloud Storage para un bucket en region-c, se te cobrarán dos saltos de red (de region-a a region-b y, luego, de region-b a region-c, o a la inversa según la dirección del conector).

Considera cuidadosamente todas las regiones involucradas cuando planifiques la ubicación de tu clúster de Connect para optimizar la latencia y el costo.

Configuración de capacidad

La configuración de capacidad requiere que establezcas la cantidad de CPU virtuales y la cantidad de memoria para cada CPU virtual de tu clúster de Connect. Puedes actualizar la capacidad de un clúster de Connect después de crearlo. A continuación, se indican las propiedades de la configuración de capacidad:

  • CPU virtuales: Es la cantidad de CPU virtuales asignadas a un clúster de Connect. El valor mínimo es 3 CPU virtuales.

  • Memoria: Es la cantidad de memoria que se asigna a cada CPU virtual. Debes aprovisionar entre 1 GiB y 8 GiB por CPU virtual. La cantidad de memoria se puede aumentar o disminuir dentro de estos límites después de crear el clúster.

    Por ejemplo, si creas un clúster con 6 CPU virtuales, la memoria mínima que puedes asignar al clúster es de 6 GiB (1 GiB por CPU virtual) y la máxima es de 48 GiB (8 GiB por CPU virtual).

La CPU virtual y la memoria asignadas a cada trabajador en un clúster de Connect tienen un impacto significativo en el rendimiento, la capacidad y el costo del clúster. A continuación, se incluye un desglose de cómo la CPU virtual y la memoria afectan a un clúster de Connect.

Recuento de CPU virtuales

  • Kafka Connect divide el trabajo de un conector en tareas. Cada tarea puede procesar datos en paralelo. Más CPU virtuales significan que se pueden ejecutar más tareas de forma simultánea, lo que genera una mayor capacidad de procesamiento.

  • Más CPU virtuales aumentan los costos de tu clúster de Connect.

Memoria

  • Kafka Connect usa memoria para almacenar datos en búfer a medida que fluyen entre los conectores y Managed Service para Apache Kafka. Una mayor memoria permite búferes más grandes. Una gran cantidad de memoria puede mejorar la capacidad de procesamiento, en especial para los flujos de datos de alto volumen. Los conectores que procesan mensajes o registros muy grandes requieren suficiente memoria para procesarlos sin generar excepciones de OutOfMemoryError.

  • Más memoria aumenta el costo de tu clúster de Connect.

  • Si usas una lógica de transformación pesada, necesitas una mayor asignación de memoria.

Tu objetivo es elegir la configuración de capacidad adecuada para tu clúster de Connect. Para ello, debes comprender la capacidad de procesamiento que puede controlar tu clúster de Connect.

Subred de trabajadores (principal)

La subred de trabajadores, también conocida como subred principal, conecta tu red de VPC al clúster de Connect. Esta subred permite que los trabajadores del clúster lleguen a los extremos de las fuentes y los receptores en la red del consumidor, como los clústeres de Managed Service para Apache Kafka o los clústeres de Kafka autohospedados.

Estos son algunos requisitos para configurar la subred del trabajador:

A los trabajadores del clúster se les asignan direcciones IP en la subred de trabajadores, con una interfaz de Private Service Connect. Los trabajadores pueden llegar a cualquier destino de red al que se pueda acceder desde la red de VPC de la subred, con los siguientes requisitos:

  • El extremo no debe estar en el rango CIDR 172.16.0.0/14. Este rango está reservado para el uso interno de Managed Service para Apache Kafka Connect.
  • Las reglas de firewall deben permitir el tráfico. Consulta Configura la seguridad para los adjuntos de red.
  • Para el tráfico de Internet, debes configurar un Cloud NAT. Por ejemplo, se requiere un Cloud NAT para que un conector de MirrorMaker replique datos de un clúster de Kafka al que se puede acceder a través de Internet.
  • Para acceder a los extremos de Private Service Connect que se encuentran en una VPC diferente de la subred del trabajador, debes asegurarte de usar una configuración de consumidor compatible (por ejemplo, NCC). Para obtener más información, consulta Información sobre el acceso a servicios publicados a través de extremos.

Dominios del DNS que se pueden resolver

Los dominios DNS que se pueden resolver, también conocidos como nombres de dominio DNS, permiten que las direcciones DNS de la red de VPC del consumidor estén disponibles para la VPC del arrendatario. Esto permite que el clúster de Connect resuelva nombres de DNS en direcciones IP, lo que facilita la comunicación con otros servicios, incluidos otros clústeres de Kafka para los conectores de MirrorMaker.

Para los dominios DNS que se pueden resolver, puedes seleccionar un clúster de Managed Service para Apache Kafka. No es necesario que configures el nombre de dominio DNS para el clúster principal de Managed Service para Apache Kafka, ya que su dirección de arranque se incluye automáticamente en la lista de dominios DNS que se pueden resolver.

Sin embargo, también puedes especificar un dominio DNS de forma manual, lo que es necesario si seleccionas un clúster de Kafka externo. El dominio DNS del clúster principal de Managed Service para Apache Kafka se incluye automáticamente. Los demás clústeres de Kafka aún requieren la configuración de dominios DNS.

Recursos de Secret Manager

Algunos conectores requieren datos sensibles, como contraseñas, como parte de su configuración. Para administrar este tipo de datos de forma segura, puedes almacenarlos en Secret Manager y otorgar acceso al secreto al clúster de Connect.

Para usar secretos de Secret Manager con Kafka Connect, haz lo siguiente:

  1. Otorga el rol de Descriptor de acceso a secretos de Secret Manager (roles/secretmanager.secretAccessor) a la cuenta de servicio de Kafka administrada. Este rol permite que tu clúster de Connect acceda a los secretos.

  2. Crea un secreto que contenga los datos sensibles. Para obtener más información, consulta Crea un secreto.

  3. Cuando crees o actualices tu clúster de Connect, especifica los secretos a los que tiene acceso el clúster. Puedes especificar hasta 32 secretos por clúster de Connect.

Los secretos se activan como archivos en los nodos de procesamiento del clúster. Los conectores tienen acceso de solo lectura a estos archivos. Cuando creas un conector, las propiedades de configuración del conector pueden hacer referencia a los secretos.

  • Para hacer referencia a la ruta de acceso a un archivo secreto, usa el siguiente formato:

    /var/secrets/PROJECT_NAME-SECRET_NAME-SECRET_VERSION
    

    Ejemplo: ssl.truststore.location=/var/secrets/project1-truststore-1

  • Para usar el valor de un secreto como valor de configuración (por ejemplo, una contraseña), usa el siguiente formato:

    ${directory:/var/secrets:PROJECT_NAME-SECRET_NAME-SECRET_VERSION}
    

    Ejemplo: password=${directory:/var/secrets:project1-database_password-3}

Reemplaza lo siguiente:

  • PROJECT_NAME: Es el nombre del proyecto de Google Cloud .
  • SECRET_NAME: Es el nombre del secreto.
  • SECRET_VERSION: Es la versión del secreto.

Etiquetas

Las etiquetas son pares clave-valor que te ayudan con la organización y la identificación. Te ayudan a organizar los clústeres de Connect. Puedes adjuntar una etiqueta a cada clúster de Connect y, luego, filtrar los recursos según sus etiquetas. Algunos ejemplos de etiquetas son environment:prod y application:web-app.

Crea un clúster de Connect

Antes de crear un clúster, revisa la documentación sobre las propiedades del clúster de Connect.

La creación de un clúster de Connect tarda entre 20 y 30 minutos.

Console

  1. En la consola de Google Cloud , ve a la página Connect Clusters.

    Ir a Connect Clusters

  2. Haga clic en Crear.

  3. En el campo Nombre del clúster de Connect, ingresa un nombre para el clúster de Connect. Para obtener más información, consulta los Lineamientos para asignarles nombres a los recursos de Managed Service para Apache Kafka.

  4. En la lista Clúster principal de Kafka, selecciona un clúster de Managed Service para Apache Kafka. Para obtener más información, consulta Clúster principal de Kafka.

  5. En la lista Región, selecciona una ubicación para el clúster de Connect. Para obtener más información sobre cómo seleccionar una ubicación, consulta Clúster principal de Kafka.

  6. En la sección Configuración de capacidad, ingresa valores para los siguientes campos o conserva los valores predeterminados.

    • En el campo vCPUs, ingresa la cantidad de CPU virtuales para el clúster.

    • En el campo Memoria, ingresa la cantidad de memoria por CPU, en GiB. El valor no puede exceder los 8 GiB por CPU.

    Para obtener más información, consulta Configuración de capacidad.

  7. En la sección Configuración de red, selecciona una red de VPC de la lista Red o deja el valor predeterminado. Esta lista se completa cuando seleccionas el clúster principal de Kafka.

  8. En la sección Subred del trabajador, selecciona una subred de la lista Subred o deja el valor predeterminado. Para obtener más información, consulta Subred de trabajadores. El campo Ruta URI de la subred se completa automáticamente cuando seleccionas la subred.

  9. Opcional: Agrega un dominio de DNS que se pueda resolver. El dominio DNS del clúster principal de Kafka se agrega automáticamente como un dominio DNS que se puede resolver. Para especificar dominios DNS adicionales, haz lo siguiente:

    1. Expande la sección Dominios del DNS que se pueden resolver.

    2. Haz clic en Agregar dominio de DNS.

    3. Para agregar el dominio DNS de un clúster existente de Managed Service para Apache Kafka, selecciona el clúster en la lista Clúster de Kafka. De lo contrario, ingresa el dominio de DNS en el campo Dominio de DNS.

    4. Haz clic en Listo.

  10. Opcional: Para agregar recursos de Secret Manager, haz lo siguiente:

    1. Expande la sección Recursos de Secret Manager.

    2. Haz clic en Agregar recurso secreto.

    3. En la lista Secreto, selecciona un secreto.

    4. En la lista Versión del secreto, selecciona una versión del secreto.

    5. Haz clic en Listo.

  11. Opcional: Agrega etiquetas para organizar tu proyecto. Para agregar una etiqueta, haz lo siguiente:

    1. Expande la sección Etiquetas.

    2. Haz clic en Agregar etiqueta.

    3. En el campo Clave, ingresa la clave de la etiqueta.

    4. En el campo Valor, ingresa el valor de la etiqueta.

  12. Haz clic en Crear.

gcloud

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

    Activa Cloud Shell

    En la parte inferior de la consola de Google Cloud , se inicia una sesión de Cloud Shell que muestra una ventana emergente con una 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. Ejecuta el comando gcloud managed-kafka connect-clusters create:

    gcloud managed-kafka connect-clusters create CONNECT_CLUSTER_ID \
        --location=LOCATION \
        --cpu=CPU \
        --memory=MEMORY \
        --primary-subnet=WORKER_SUBNET \
        --kafka-cluster=KAFKA_CLUSTER \
        [--project=PROJECT_ID] \
        [--secret=SECRET] \
        [--dns-name=DNS_DOMAIN_NAME] \
        [--config-file=CONFIG_FILE] \
        [--labels=LABELS]
        [--async]
    

    Reemplaza lo siguiente:

    • CONNECT_CLUSTER_ID: Es el ID o el nombre del clúster de Connect. Si necesitas ayuda para asignarle un nombre a un clúster de Connect, consulta los Lineamientos para asignarles nombres a los recursos de Managed Service para Apache Kafka. El nombre de un clúster de Connect es inmutable.

    • LOCATION: Es la ubicación en la que creas el clúster de Connect. Debe ser una región Google Cloudadmitida. No puedes cambiar la ubicación de un clúster de Connect después de crearlo. Para ver una lista de las ubicaciones disponibles, consulta las ubicaciones de Managed Service para Apache Kafka. Para obtener más información sobre las recomendaciones de ubicación, consulta Clúster principal de Kafka.

    • CPU: Es la cantidad de CPU virtuales para el clúster de Connect. El valor mínimo es 3 CPU virtuales. Consulta Recuento de CPU virtuales.

    • MEMORY: Es la cantidad de memoria para el clúster de Connect. Usa las unidades "MB", "MiB", "GB", "GiB", "TB" o "TiB". Por ejemplo, "3 GiB". Debes aprovisionar entre 1 GiB y 8 GiB por CPU virtual. Consulta Memory.

    • WORKER_SUBNET: Es la subred de trabajadores del clúster de Connect.

      El formato de la subred es projects/PROJECT_ID/regions/REGION/subnetworks/SUBNET_ID.

      La subred de trabajadores debe estar en la misma región que el clúster de Connect.

    • PROJECT_ID: (Opcional) ID del proyecto deGoogle Cloud . Si no se proporciona, se usa el proyecto actual.

    • KAFKA_CLUSTER: Es el ID o el nombre completamente calificado del clúster principal de Managed Service para Apache Kafka asociado con el clúster de Connect. Consulta Clúster de Kafka. El formato del clúster de Kafka es projects/PROJECT_ID/locations/LOCATION/clusters/CLUSTER_ID.

      No puedes actualizar a un clúster de Kafka diferente después de crear el clúster de Connect.

    • SECRET: (Opcional) Son los secretos que se cargarán en los trabajadores. Se deben proporcionar las versiones exactas de los secretos de Secret Manager. Los alias no son compatibles. Se pueden cargar hasta 32 secretos en un clúster. Formato: projects/PROJECT_ID/secrets/SECRET_NAME/versions/VERSION_ID

    • DNS_DOMAIN_NAME: (Opcional) Nombres de dominio DNS de la subred que se harán visibles para el clúster de Connect. El clúster de Connect puede acceder a los recursos con nombres de dominio en lugar de depender de direcciones IP. Consulta Intercambio de tráfico de DNS.

    • LABELS: Son etiquetas opcionales para asociar con el clúster. Para obtener más información sobre el formato de las etiquetas, consulta Etiquetas. Lista de pares clave-valor de etiquetas que se agregarán. Las claves deben comenzar con un carácter en minúscula y contener solo guiones (-), guiones bajos (_), caracteres en minúscula y números. Los valores deben contener solo guiones (-), guiones bajos (_), caracteres en minúscula y números.

    • CONFIG_FILE: (Opcional) Es la ruta al archivo JSON o YAML que contiene la configuración que se anula de los valores predeterminados del clúster o del conector. Este archivo también admite JSON o YAML intercalados.

    • --async: (Opcional) Se muestra de inmediato, sin necesidad de esperar a que se complete la operación en curso. Con la marca --async, puedes continuar con otras tareas mientras se crea el clúster en segundo plano. Si no usas la marca, el sistema espera a que se complete la operación antes de devolver una respuesta. Debes esperar a que el clúster se actualice por completo antes de continuar con otras tareas.

    Obtendrás una respuesta similar a la siguiente:

    Create request issued for: [sample-connectcluster]
    Check operation [projects/test-project/locations/us-east1/operations/operation-1753590328249-63ae19098cc06-64300a0a-06512d02] for status.
    

    Almacena el OPERATION_ID para hacer un seguimiento del progreso. Por ejemplo, el valor aquí es operation-1753590328249-63ae19098cc06-64300a0a-06512d02.

Terraform

Puedes usar un recurso de Terraform para crear un clúster de Connect.

resource "google_managed_kafka_connect_cluster" "default" {
  provider           = google-beta
  project            = data.google_project.default.project_id
  connect_cluster_id = "my-connect-cluster-id"
  location           = "us-central1"
  kafka_cluster      = google_managed_kafka_cluster.default.id
  capacity_config {
    vcpu_count   = 12
    memory_bytes = 12884901888 # 12 GiB
  }
  gcp_config {
    access_config {
      network_configs {
        primary_subnet = google_compute_subnetwork.default.id
      }
    }
  }
}

Si deseas obtener más información para aplicar o quitar una configuración de Terraform, consulta los comandos básicos de Terraform.

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 Managed Service para Apache Kafka en Go.

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"

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

func createConnectCluster(w io.Writer, projectID, region, clusterID, kafkaCluster string, opts ...option.ClientOption) error {
	// projectID := "my-project-id"
	// region := "us-central1"
	// clusterID := "my-connect-cluster"
	// kafkaCluster := "projects/my-project-id/locations/us-central1/clusters/my-kafka-cluster"
	ctx := context.Background()
	client, err := managedkafka.NewManagedKafkaConnectClient(ctx, opts...)
	if err != nil {
		return fmt.Errorf("managedkafka.NewManagedKafkaConnectClient got err: %w", err)
	}
	defer client.Close()

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

	// Capacity configuration with 12 vCPU and 12 GiB RAM
	capacityConfig := &managedkafkapb.CapacityConfig{
		VcpuCount:   12,
		MemoryBytes: 12884901888, // 12 GiB in bytes
	}

	// Optionally, you can also specify accessible subnets and resolvable DNS
	// domains as part of your network configuration. For example:
	// networkConfigs := []*managedkafkapb.ConnectNetworkConfig{
	// 	{
	// 		PrimarySubnet:      primarySubnet,
	// 		AdditionalSubnets:  []string{"subnet-1", "subnet-2"},
	// 		DnsDomainNames:     []string{"domain-1", "domain-2"},
	// 	},
	// }

	connectCluster := &managedkafkapb.ConnectCluster{
		Name:           clusterPath,
		KafkaCluster:   kafkaCluster,
		CapacityConfig: capacityConfig,
	}

	req := &managedkafkapb.CreateConnectClusterRequest{
		Parent:           locationPath,
		ConnectClusterId: clusterID,
		ConnectCluster:   connectCluster,
	}
	op, err := client.CreateConnectCluster(ctx, req)
	if err != nil {
		return fmt.Errorf("client.CreateConnectCluster got err: %w", err)
	}
	// The duration of this operation can vary considerably, typically taking 5-15 minutes.
	resp, err := op.Wait(ctx)
	if err != nil {
		return fmt.Errorf("op.Wait got err: %w", err)
	}
	fmt.Fprintf(w, "Created connect cluster: %s\n", resp.Name)
	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.RetryingFuture;
import com.google.api.gax.retrying.TimedRetryAlgorithm;
import com.google.cloud.managedkafka.v1.CapacityConfig;
import com.google.cloud.managedkafka.v1.ConnectAccessConfig;
import com.google.cloud.managedkafka.v1.ConnectCluster;
import com.google.cloud.managedkafka.v1.ConnectGcpConfig;
import com.google.cloud.managedkafka.v1.ConnectNetworkConfig;
import com.google.cloud.managedkafka.v1.CreateConnectClusterRequest;
import com.google.cloud.managedkafka.v1.LocationName;
import com.google.cloud.managedkafka.v1.ManagedKafkaConnectClient;
import com.google.cloud.managedkafka.v1.ManagedKafkaConnectSettings;
import com.google.cloud.managedkafka.v1.OperationMetadata;
import java.time.Duration;
import java.util.concurrent.ExecutionException;

public class CreateConnectCluster {

  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-connect-cluster";
    String subnet = "my-subnet"; // e.g. projects/my-project/regions/my-region/subnetworks/my-subnet
    String kafkaCluster = "my-kafka-cluster"; // The Kafka cluster to connect to
    int cpu = 12;
    long memoryBytes = 12884901888L; // 12 GiB
    createConnectCluster(projectId, region, clusterId, subnet, kafkaCluster, cpu, memoryBytes);
  }

  public static void createConnectCluster(
      String projectId,
      String region,
      String clusterId,
      String subnet,
      String kafkaCluster,
      int cpu,
      long memoryBytes)
      throws Exception {
    CapacityConfig capacityConfig = CapacityConfig.newBuilder().setVcpuCount(cpu)
        .setMemoryBytes(memoryBytes).build();
    ConnectNetworkConfig networkConfig = ConnectNetworkConfig.newBuilder()
        .setPrimarySubnet(subnet)
        .build();
    // Optionally, you can also specify additional accessible subnets and resolvable
    // DNS domains as part of your network configuration. For example:
    // .addAllAdditionalSubnets(List.of("subnet-1", "subnet-2"))
    // .addAllDnsDomainNames(List.of("dns-1", "dns-2"))
    ConnectGcpConfig gcpConfig = ConnectGcpConfig.newBuilder()
        .setAccessConfig(ConnectAccessConfig.newBuilder().addNetworkConfigs(networkConfig).build())
        .build();
    ConnectCluster connectCluster = ConnectCluster.newBuilder()
        .setCapacityConfig(capacityConfig)
        .setGcpConfig(gcpConfig)
        .setKafkaCluster(kafkaCluster)
        .build();

    // Create the settings to configure the timeout for polling operations
    ManagedKafkaConnectSettings.Builder settingsBuilder = ManagedKafkaConnectSettings.newBuilder();
    TimedRetryAlgorithm timedRetryAlgorithm = OperationTimedPollAlgorithm.create(
        RetrySettings.newBuilder()
            .setTotalTimeoutDuration(Duration.ofHours(1L))
            .build());
    settingsBuilder.createConnectClusterOperationSettings()
        .setPollingAlgorithm(timedRetryAlgorithm);

    try (ManagedKafkaConnectClient managedKafkaConnectClient = ManagedKafkaConnectClient
        .create(settingsBuilder.build())) {
      CreateConnectClusterRequest request = CreateConnectClusterRequest.newBuilder()
          .setParent(LocationName.of(projectId, region).toString())
          .setConnectClusterId(clusterId)
          .setConnectCluster(connectCluster)
          .build();

      // The duration of this operation can vary considerably, typically taking
      // between 10-30 minutes.
      OperationFuture<ConnectCluster, OperationMetadata> future = managedKafkaConnectClient
          .createConnectClusterOperationCallable().futureCall(request);

      // Get the initial LRO and print details.
      OperationSnapshot operation = future.getInitialFuture().get();
      System.out.printf(
          "Connect 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)
      ConnectCluster response = future.get();
      System.out.printf("Created connect cluster: %s\n", response.getName());
    } catch (ExecutionException e) {
      System.err.printf("managedKafkaConnectClient.createConnectCluster got err: %s\n", 
          e.getMessage());
      throw e;
    }
  }
}

Python

Antes de probar este ejemplo, sigue las instrucciones de configuración de Python en Instala las bibliotecas cliente. Si deseas 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.cloud.managedkafka_v1.services.managed_kafka_connect import ManagedKafkaConnectClient
from google.cloud.managedkafka_v1.types import ConnectCluster, CreateConnectClusterRequest, ConnectNetworkConfig

# TODO(developer): Update with your values.
# project_id = "my-project-id"
# region = "us-central1"
# connect_cluster_id = "my-connect-cluster"
# kafka_cluster_id = "my-kafka-cluster"
# primary_subnet = "projects/my-project-id/regions/us-central1/subnetworks/default"
# cpu = 12
# memory_bytes = 12884901888  # 12 GiB

connect_client = ManagedKafkaConnectClient()
kafka_client = managedkafka_v1.ManagedKafkaClient()

parent = connect_client.common_location_path(project_id, region)
kafka_cluster_path = kafka_client.cluster_path(project_id, region, kafka_cluster_id)

connect_cluster = ConnectCluster()
connect_cluster.name = connect_client.connect_cluster_path(project_id, region, connect_cluster_id)
connect_cluster.kafka_cluster = kafka_cluster_path
connect_cluster.capacity_config.vcpu_count = cpu
connect_cluster.capacity_config.memory_bytes = memory_bytes
connect_cluster.gcp_config.access_config.network_configs = [ConnectNetworkConfig(primary_subnet=primary_subnet)]
# Optionally, you can also specify accessible subnets and resolvable DNS domains as part of your network configuration.
# For example:
# connect_cluster.gcp_config.access_config.network_configs = [
#     ConnectNetworkConfig(
#         primary_subnet=primary_subnet,
#         additional_subnets=additional_subnets,
#         dns_domain_names=dns_domain_names,
#     )
# ]

request = CreateConnectClusterRequest(
    parent=parent,
    connect_cluster_id=connect_cluster_id,
    connect_cluster=connect_cluster,
)

try:
    operation = connect_client.create_connect_cluster(request=request)
    print(f"Waiting for operation {operation.operation.name} to complete...")
    # Creating a Connect cluster can take 10-40 minutes.
    response = operation.result(timeout=3000)
    print("Created Connect cluster:", response)
except GoogleAPICallError as e:
    print(f"The operation failed with error: {e}")

Supervisa la operación de creación del clúster

Solo puedes ejecutar el siguiente comando si ejecutaste la gcloud CLI para crear el clúster de Connect.

  • Por lo general, la creación de un clúster de Connect tarda entre 20 y 30 minutos. Para hacer un seguimiento del progreso de la creación del clúster, el comando gcloud managed-kafka connect-clusters create usa una operación de larga duración (LRO), que puedes supervisar con el siguiente comando:

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

    Reemplaza lo siguiente:

    • OPERATION_ID con el valor del ID de operación de la sección anterior
    • LOCATION con el valor de la ubicación de la sección anterior

Próximos pasos

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