删除 Connect 集群

删除 Connect 集群会擦除所有关联的数据,包括存储在其主 Kafka 集群中的连接器配置。此操作无法撤销。

如需删除 Connect 集群,您可以使用 Google Cloud 控制台、 gcloud CLI、客户端库或 Managed Kafka API。您无法使用开源 Apache Kafka API 删除 Connect 集群。

删除 Connect 集群所需的角色和权限

如需获得删除 Connect 集群所需的权限,请让您的管理员为您授予项目的Managed Kafka Connect Cluster Editor (roles/managedkafka.connectClusterEditor) IAM 角色。如需详细了解如何授予角色,请参阅管理对项目、文件夹和组织的访问权限

此预定义角色包含 删除 Connect 集群所需的权限。如需查看所需的确切权限,请展开所需权限部分:

所需权限

您需要具备以下权限才能删除 Connect 集群:

  • 授予对 Connect 集群的删除 Connect 集群权限: managedkafka.connectClusters.delete
  • 授予对指定位置的列出 Connect 集群权限。只有使用 Google Cloud 控制台删除 Connect 集群时,才需要此权限: managedkafka.connectClusters.list

您也可以使用自定义角色或其他预定义角色来获取这些权限。

删除 Connect 集群

  • 了解数据丢失的影响: 删除 Connect 集群会擦除存储在 Connect 集群本身内的所有数据。 这包括以下内容:

    • 连接器及其配置

    • 由 Connect 集群直接管理的任何其他数据

    删除 Connect 集群不会删除源 Kafka 集群或目标 Kafka 集群中的数据。 如果您使用源连接器将数据移至 Kafka 主题,则删除 Connect 集群不会删除已发布到该 Kafka 主题的数据。同样,删除 Connect 集群也不会删除与 Connect 集群关联的 Kafka 集群。

  • 规划服务中断: 任何依赖于 Connect 集群读取或写入数据的应用或服务都可能会中断。 请在删除集群之前规划此服务中断。

  • 查看结算方面的影响: 删除集群后,您将停止产生集群费用。您可能仍需支付删除之前使用的资源费用。

  • 预期异步操作: 默认情况下,集群删除是异步的。该命令会立即返回,您可以单独跟踪删除进度。

控制台

  1. 在 Google Cloud 控制台中,前往 Connect 集群 页面。

    前往 Connect 集群

  2. 选择要删除的 Connect 集群。您可以选择多个。

  3. 点击删除

gcloud

  1. 在 Google Cloud 控制台中,激活 Cloud Shell。

    激活 Cloud Shell

    Cloud Shell 会话随即会在控制台的底部启动,并显示命令行提示符。 Google Cloud Cloud Shell 是一个已安装 Google Cloud CLI 且已为当前项目设置值的 Shell 环境。该会话可能需要几秒钟来完成初始化。

  2. 使用 gcloud managed-kafka connect-clusters delete 命令删除 Connect 集群:

    gcloud managed-kafka connect-clusters delete CONNECT_CLUSTER \
        --location=LOCATION [--async]
    

    替换以下内容:

    • CONNECT_CLUSTER:要删除的 Connect 集群的 ID 。
    • LOCATION:Connect 集群的位置。

    以下标志是可选的:

    • --async:立即返回结果,而无需等待正在进行的操作完成。

Go

在试用此示例之前,请按照 安装客户端库中的 Go 设置说明进行操作。如需了解详情, 请参阅 Managed Service for Apache Kafka Go API 参考文档

如需向 Managed Service for Apache Kafka 进行身份验证,请设置应用默认凭据(ADC)。 如需了解详情, 请参阅 为本地开发环境设置 ADC

import (
	"context"
	"fmt"
	"io"

	"cloud.google.com/go/managedkafka/apiv1/managedkafkapb"
	"google.golang.org/api/option"

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

func deleteConnectCluster(w io.Writer, projectID, region, clusterID string, opts ...option.ClientOption) error {
	// projectID := "my-project-id"
	// region := "us-central1"
	// clusterID := "my-connect-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()

	clusterPath := fmt.Sprintf("projects/%s/locations/%s/connectClusters/%s", projectID, region, clusterID)
	req := &managedkafkapb.DeleteConnectClusterRequest{
		Name: clusterPath,
	}
	op, err := client.DeleteConnectCluster(ctx, req)
	if err != nil {
		return fmt.Errorf("client.DeleteConnectCluster got err: %w", err)
	}
	err = op.Wait(ctx)
	if err != nil {
		return fmt.Errorf("op.Wait got err: %w", err)
	}
	fmt.Fprint(w, "Deleted connect cluster\n")
	return nil
}

Java

在试用此示例之前,请按照 安装客户端库中的 Java 设置说明进行操作。如需了解详情, 请参阅 Managed Service for Apache Kafka Java API 参考文档

如需向 Managed Service for Apache Kafka 进行身份验证,请设置应用默认凭据。 如需了解详情,请参阅 为本地开发环境设置 ADC


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.api.gax.rpc.ApiException;
import com.google.cloud.managedkafka.v1.ConnectClusterName;
import com.google.cloud.managedkafka.v1.DeleteConnectClusterRequest;
import com.google.cloud.managedkafka.v1.ManagedKafkaConnectClient;
import com.google.cloud.managedkafka.v1.ManagedKafkaConnectSettings;
import com.google.cloud.managedkafka.v1.OperationMetadata;
import com.google.protobuf.Empty;
import java.io.IOException;
import java.time.Duration;

public class DeleteConnectCluster {

  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";
    deleteConnectCluster(projectId, region, clusterId);
  }

  public static void deleteConnectCluster(String projectId, String region, String clusterId)
      throws Exception {

    // 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.deleteConnectClusterOperationSettings()
        .setPollingAlgorithm(timedRetryAlgorithm);

    try (ManagedKafkaConnectClient managedKafkaConnectClient = ManagedKafkaConnectClient.create(
        settingsBuilder.build())) {
      DeleteConnectClusterRequest request = DeleteConnectClusterRequest.newBuilder()
          .setName(ConnectClusterName.of(projectId, region, clusterId).toString())
          .build();
      OperationFuture<Empty, OperationMetadata> future = managedKafkaConnectClient
          .deleteConnectClusterOperationCallable().futureCall(request);

      // Get the initial LRO and print details. CreateConnectCluster contains sample
      // code for polling logs.
      OperationSnapshot operation = future.getInitialFuture().get();
      System.out.printf(
          "Connect cluster deletion started. Operation name: %s\nDone: %s\nMetadata: %s\n",
          operation.getName(),
          operation.isDone(),
          future.getMetadata().get().toString());

      future.get();
      System.out.println("Deleted connect cluster");
    } catch (IOException | ApiException e) {
      System.err.printf("managedKafkaConnectClient.deleteConnectCluster got err: %s\n", 
          e.getMessage());
    }
  }
}

Python

在试用此示例之前,请按照 安装客户端库中的 Python 设置说明进行操作。如需了解详情, 请参阅 Managed Service for Apache Kafka Python API 参考文档

如需向 Managed Service for Apache Kafka 进行身份验证,请设置应用默认凭据。 如需了解详情, 请参阅为本地开发环境设置 ADC

from google.api_core.exceptions import GoogleAPICallError
from google.cloud.managedkafka_v1.services.managed_kafka_connect import (
    ManagedKafkaConnectClient,
)
from google.cloud import managedkafka_v1

# TODO(developer)
# project_id = "my-project-id"
# region = "us-central1"
# connect_cluster_id = "my-connect-cluster"

connect_client = ManagedKafkaConnectClient()

request = managedkafka_v1.DeleteConnectClusterRequest(
    name=connect_client.connect_cluster_path(project_id, region, connect_cluster_id),
)

try:
    operation = connect_client.delete_connect_cluster(request=request)
    print(f"Waiting for operation {operation.operation.name} to complete...")
    operation.result()
    print("Deleted Connect cluster")
except GoogleAPICallError as e:
    print(f"The operation failed with error: {e}")

接下来怎么做?

Apache Kafka® 是 Apache Software Foundation 或其关联公司在美国和/或其他国家/地区的注册 商标。