列出所有连接器

列出 Connect 集群中运行的所有连接器,可以概览已配置的数据集成。您可以监控连接器的运行状况和状态,找出潜在问题,并有效地管理数据流。

如需列出 Connect 集群中的所有连接器,您可以使用 Google Cloud 控制台、 gcloud CLI、Managed Service for Apache Kafka 客户端库或 Managed Kafka API。您无法使用开源 Apache Kafka API 列出连接器。

列出所有连接器所需的角色和权限

如需获取列出所有连接器所需的权限,请让管理员在您的项目上向您授予Managed Kafka 查看者 (roles/managedkafka.viewer) IAM 角色。如需详细了解如何授予角色,请参阅管理对项目、文件夹和组织的访问权限

此预定义角色包含 列出所有连接器所需的权限。如需查看所需的确切权限,请展开所需权限部分:

所需权限

如需列出所有连接器,您需要以下权限:

  • 在父级 Connect 集群上授予列出连接器权限: managedkafka.connectors.list
  • 在父级 Connect 集群上授予获取连接器详情权限: managedkafka.connectors.get

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

查看所有连接器

此视图提供了一种快速监控连接器状态的方法,并可找出需要注意的连接器。然后,您可以根据需要深入了解各个连接器,以查看其详细信息和配置。

控制台

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

    前往 Connect 集群

  2. 点击要列出连接器的 Connect 集群。

    系统会显示 Connect 集群详情 页面。

  3. 资源 标签页会显示集群中运行的所有连接器的列表。该列表包含每个连接器的以下信息:

    • 名称:连接器的名称。
    • 状态:连接器的操作状态。例如,正在运行、失败。
    • 连接器类型:连接器插件的类型。

    您可以使用过滤条件 选项按名称搜索特定连接器。

gcloud

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

    激活 Cloud Shell

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

  2. 使用 gcloud managed-kafka connectors list 命令列出连接器:

    gcloud managed-kafka connectors list CONNECT_CLUSTER_ID \
        --location=LOCATION
    
  3. 如需进一步优化连接器列表,您可以使用其他标志:

    gcloud managed-kafka connectors list CONNECT_CLUSTER_ID \
        --location=LOCATION \
        [--filter=EXPRESSION] \
        [--limit=LIMIT] \
        [--page-size=PAGE_SIZE] \
        [--sort-by=SORT_BY]
    

    替换以下内容:

    • CONNECT_CLUSTER_ID:必填。包含您要列出的连接器的 Connect 集群的 ID。
    • LOCATION:必填。包含您要列出的连接器的 Connect 集群 的位置。
    • EXPRESSION:(可选)要应用于列表的布尔过滤 条件表达式。如果表达式的计算结果为 True,则该项会包含在列表中。如需了解更多详情和示例,请运行 gcloud topic filters

      示例:

      • 如需仅列出处于“RUNNING”状态的连接器:

        --filter="state=RUNNING"
        
      • 如需仅列出“Pub/Sub Sink”连接器:

        --filter="connector_plugin='Pub/Sub Sink'"
        
      • 如需列出名称包含“prod”的连接器:

        --filter="name ~ 'prod'"
        
      • 如需列出“FAILED”或“Pub/Sub Source”插件的连接器:

        --filter="state=FAILED OR connector_plugin='Pub/Sub Source'"
        
    • LIMIT:(可选)要显示的最大 连接器数量。如果未指定,系统会列出所有连接器。

    • PAGE_SIZE:(可选)每页显示的结果数。如果未指定,服务会确定合适的页面大小。

    • SORT_BY:(可选)以逗号分隔列表形式指定的字段,用于排序。默认的排序顺序是升序。如需按降序排序,请在字段前面加上 ~。支持的字段可能是 namestate

包含排序的示例命令:

gcloud managed-kafka connectors list test-connect-cluster \
    --location=us-central1 \
    --sort-by=~state,name

包含过滤和限制的示例命令:

gcloud managed-kafka connectors list test-connect-cluster \
    --location=us-central1 \
    --filter="state=RUNNING AND connector_plugin='Pub/Sub Sink'" \
    --limit=5

输出示例:

NAME                                    STATE     CONNECTOR_PLUGIN
pubsub-sink-connector                   RUNNING   Pub/Sub Sink
another-pubsub-sink                     RUNNING   Pub/Sub Sink
prod-pubsub-sink                        RUNNING   Pub/Sub Sink

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/iterator"
	"google.golang.org/api/option"

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

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

	parent := fmt.Sprintf("projects/%s/locations/%s/connectClusters/%s", projectID, region, connectClusterID)
	req := &managedkafkapb.ListConnectorsRequest{
		Parent: parent,
	}
	connectorIter := client.ListConnectors(ctx, req)
	for {
		connector, err := connectorIter.Next()
		if err == iterator.Done {
			break
		}
		if err != nil {
			return fmt.Errorf("connectorIter.Next() got err: %w", err)
		}
		fmt.Fprintf(w, "Got connector: %v", connector)
	}
	return nil
}

Java

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

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

import com.google.api.gax.rpc.ApiException;
import com.google.cloud.managedkafka.v1.ConnectClusterName;
import com.google.cloud.managedkafka.v1.Connector;
import com.google.cloud.managedkafka.v1.ManagedKafkaConnectClient;
import java.io.IOException;

public class ListConnectors {

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

  public static void listConnectors(String projectId, String region, String clusterId)
      throws IOException {
    try (ManagedKafkaConnectClient managedKafkaConnectClient = ManagedKafkaConnectClient.create()) {
      ConnectClusterName parent = ConnectClusterName.of(projectId, region, clusterId);
      // This operation is handled synchronously.
      for (Connector connector : managedKafkaConnectClient.listConnectors(parent).iterateAll()) {
        System.out.println(connector.getAllFields());
      }
    } catch (IOException | ApiException e) {
      System.err.printf("managedKafkaConnectClient.listConnectors got err: %s\n", e.getMessage());
    }
  }
}

Python

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

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

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

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

connect_client = ManagedKafkaConnectClient()

request = managedkafka_v1.ListConnectorsRequest(
    parent=connect_client.connect_cluster_path(project_id, region, connect_cluster_id),
)

try:
    response = connect_client.list_connectors(request=request)
    for connector in response:
        print("Got connector:", connector)
except GoogleAPICallError as e:
    print(f"Failed to list connectors with error: {e}")

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