创建 Cloud Storage Sink 连接器

Cloud Storage 接收器连接器会将数据从 Kafka 流式传输到 Cloud Storage 存储分区,让您能够以经济实惠且伸缩的方式存储和处理大量数据。

Cloud Storage 接收器连接器的使用场景包括:

  • 数据湖注入。将 Kafka 数据存储在数据湖中,以便进行长期归档和批量处理。

  • 归档数据以满足监管要求。

准备工作

在创建 Cloud Storage 接收器连接器之前,请确保您已具备以下条件:

所需的角色和权限

如需获得创建连接器所需的权限,请让您的管理员为您授予项目的Managed Kafka Connector Editor (roles/managedkafka.connectorEditor) IAM 角色。如需详细了解如何授予角色,请参阅管理对项目、文件夹和组织的访问权限

此预定义角色可提供创建连接器所需的权限。如需查看所需的确切权限,请展开所需权限部分:

所需权限

创建连接器需要以下权限:

  • 创建连接器: managedkafka.connectors.create

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

授予写入 Cloud Storage 存储桶的权限

Managed Kafka 服务帐号必须具有以下权限才能将消息写入 Cloud Storage 存储桶:

  • storage.objects.create
  • storage.objects.delete

在包含 Cloud Storage 存储桶的项目中,向服务帐号授予 Storage Object User (roles/storage.objectUser) 角色。

Managed Kafka 服务帐号采用以下格式: service-PROJECT_NUMBER@gcp-sa-managedkafka.iam.gserviceaccount.com,其中 PROJECT_NUMBER 是 Connect 集群的项目编号。

如果您的 Connect 集群与 Managed Service for Apache Kafka 集群位于不同的项目中,请参阅 在不同的项目中创建 Connect 集群

Cloud Storage 接收器连接器的工作原理

Cloud Storage 接收器连接器会从一个或多个 Kafka 主题提取数据,并将这些数据写入单个 Cloud Storage 存储桶中的对象。

以下是 Cloud Storage 接收器连接器复制数据的详细说明:

  • 连接器会使用源集群中的一个或多个 Kafka 主题中的消息。

  • 连接器会将数据写入您在连接器配置中指定的目标 Cloud Storage 存储桶。

  • 连接器会在将数据写入 Cloud Storage 存储桶时,通过引用连接器配置中的特定属性来格式化数据。 默认情况下,输出文件采用 CSV 格式。您可以配置 format.output.type 属性以指定不同的输出格式,例如 JSON。

  • 连接器还会为写入 Cloud Storage 存储桶的文件命名。您可以使用 file.name.prefixfile.name.template 属性自定义文件名。例如,您可以在文件名中添加 Kafka 主题名称或消息键。

  • Kafka 记录包含三个组成部分:标头、键和值。

    • 您可以通过设置 format.output.fields 来在输出文件中添加标头,以包含标头。 例如,format.output.fields=value,headers

    • 您可以通过设置 format.output.fields 来在输出文件中添加键,以包含 key。例如,format.output.fields=key,value,headers

      您还可以通过在 file.name.template 属性中添加 key,使用键对记录进行分组。

  • 默认情况下,您可以在输出文件中添加值,因为 format.output.fields 默认设置为 value

  • 连接器会将转换后的格式化数据写入指定的 Cloud Storage 存储桶。

  • 如果您使用 file.compression.type 属性配置文件压缩,连接器会压缩存储在 Cloud Storage 存储桶中的文件。

  • 转换器配置受 format.output.type 属性的限制。

    • 例如,当 format.output.type 设置为 csv 时,键转换器必须为 org.apache.kafka.connect.converters.ByteArrayConverterorg.apache.kafka.connect.storage.StringConverter,值转换器必须为 org.apache.kafka.connect.converters.ByteArrayConverter

    • format.output.type 设置为 json 时,即使 value.converter.schemas.enable 属性为 true,值和键架构也不会与输出文件中的数据一起写入。

  • tasks.max 属性控制连接器的并行级别。增加 tasks.max 可以提高吞吐量,但实际并行度受 Kafka 主题中的分区数量的限制。

Cloud Storage 接收器连接器的属性

创建 Cloud Storage 接收器连接器时,请指定以下属性。

连接器名称

连接器的名称或 ID。如需了解有关如何命名资源的指南, 请参阅 Managed Service for Apache Kafka 资源命名指南。 名称是不可变的。

连接器插件类型

在 Google Cloud 控制台中,选择 Cloud Storage 接收器 作为连接器插件类型。如果您不使用界面来配置连接器,还必须指定连接器类。

主题

连接器从中提取消息的 Kafka 主题。 您可以指定一个或多个主题,也可以使用正则表达式来匹配多个主题。例如,topic.* 可匹配所有以“topic”开头的主题。这些主题必须存在于与 Connect 集群关联的 Managed Service for Apache Kafka 集群中。

Cloud Storage 存储桶

选择或创建用于存储数据的 Cloud Storage 存储桶。

配置

您可以在此部分中为 Cloud Storage 接收器连接器指定其他特定于连接器的配置属性。

由于 Kafka 主题中的数据可以采用各种格式(例如 Avro、JSON 或原始字节),因此配置的关键部分涉及指定转换器。 转换器会将数据从 Kafka 主题中使用的格式转换为 Kafka Connect 的标准化内部格式。然后,Cloud Storage 接收器连接器会获取此内部数据,并将其转换为 Cloud Storage 存储桶所需的格式,然后再写入。

如需了解有关转换器在 Kafka Connect 中的作用、 支持的转换器类型和常见配置选项的更多一般信息, 请参阅 转换器

以下是一些特定于 Cloud Storage 接收器连接器的配置:

  • gcs.credentials.default:是否自动从执行环境中发现凭据。 Google Cloud 必须设置为 true

  • gcs.bucket.name:指定写入数据的 Cloud Storage 存储桶的名称。必须设置。

  • file.compression.type:设置存储在 Cloud Storage 存储桶中的文件的压缩类型。示例包括 gzipsnappyzstdnone。默认值为 none

  • file.name.prefix:要添加到存储在 Cloud Storage 存储桶中的每个文件名称的前缀。默认值为空。

  • format.output.type:用于将数据写入 Cloud Storage 输出文件的数据格式类型。支持的值包括: csvjsonjsonlparquet。默认值为 csv

如需查看特定于此连接器的可用配置属性的列表,请参阅 Cloud Storage 接收器连接器配置

创建 Cloud Storage 接收器连接器

在创建连接器之前,请查看 Cloud Storage 接收器连接器的属性文档

控制台

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

    前往 Connect 集群

  2. 点击要为其创建连接器的 Connect 集群。

    此时会显示 Connect 集群详情 页面。

  3. 点击创建连接器

    此时会显示创建 Kafka 连接器 页面。

  4. 对于连接器名称,请输入一个字符串。

    如需了解有关如何命名连接器的指南,请参阅 Managed Service for Apache Kafka 资源命名指南

  5. 对于连接器插件,请选择Cloud Storage 接收器

  6. 指定可从中流式传输数据的主题

  7. 选择用于存储数据的存储分区

  8. (可选)在配置 部分中配置其他设置。

  9. 选择任务重启政策 。如需了解详情,请参阅 任务重启政策

  10. 点击创建

gcloud

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

    激活 Cloud Shell

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

  2. 运行 gcloud managed-kafka connectors create 命令:

    gcloud managed-kafka connectors create CONNECTOR_ID \
        --location=LOCATION \
        --connect-cluster=CONNECT_CLUSTER_ID \
        --config-file=CONFIG_FILE
    

    替换以下内容:

    • CONNECTOR_ID:连接器的 ID 或名称。 如需了解有关如何命名连接器的指南,请参阅 Managed Service for Apache Kafka 资源命名指南。 连接器的名称是不可变的。

    • LOCATION:您在其中创建 连接器的位置。此位置必须与您创建 Connect 集群的位置相同。

    • CONNECT_CLUSTER_ID:在其中创建连接器的 Connect 集群的 ID。

    • CONFIG_FILE:连接器的 YAML 配置文件 的路径。

    以下是 Cloud Storage 接收器连接器的配置文件示例:

    connector.class: "io.aiven.kafka.connect.gcs.GcsSinkConnector"
    tasks.max: "1"
    topics: "GMK_TOPIC_ID"
    gcs.bucket.name: "GCS_BUCKET_NAME"
    gcs.credentials.default: "true"
    format.output.type: "json"
    name: "GCS_SINK_CONNECTOR_ID"
    value.converter: "org.apache.kafka.connect.json.JsonConverter"
    value.converter.schemas.enable: "false"
    key.converter: "org.apache.kafka.connect.storage.StringConverter"
    

    替换以下内容:

    • GMK_TOPIC_ID:Managed Service for Apache Kafka 主题的 ID,数据会从该主题流向 Cloud Storage 接收器连接器。

    • GCS_BUCKET_NAME:充当流水线接收器的 Cloud Storage 存储桶的名称。

    • GCS_SINK_CONNECTOR_ID:Cloud Storage 接收器连接器的 ID 或名称。如需了解有关如何命名 连接器的指南,请参阅 Managed Service for Apache Kafka 资源命名指南。 连接器的名称是不可变的。

Terraform

您可以使用 Terraform 资源创建 连接器

resource "google_managed_kafka_connector" "example-cloud-storage-sink-connector" {
  project         = data.google_project.default.project_id
  connector_id    = "my-gcs-sink-connector"
  connect_cluster = google_managed_kafka_connect_cluster.default.connect_cluster_id
  location        = "us-central1"

  configs = {
    "connector.class"                = "io.aiven.kafka.connect.gcs.GcsSinkConnector"
    "tasks.max"                      = "3"
    "topics"                         = "GMK_TOPIC_ID"
    "gcs.bucket.name"                = "GCS_BUCKET_NAME"
    "gcs.credentials.default"        = "true"
    "format.output.type"             = "json"
    "name"                           = "my-gcs-sink-connector"
    "value.converter"                = "org.apache.kafka.connect.json.JsonConverter"
    "value.converter.schemas.enable" = "false"
    "key.converter"                  = "org.apache.kafka.connect.storage.StringConverter"
  }
  provider = google-beta
}

如需了解如何应用或移除 Terraform 配置,请参阅 基本 Terraform 命令

Go

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

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

import (
	"context"
	"fmt"
	"io"

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

// createCloudStorageSinkConnector creates a Cloud Storage Sink connector.
func createCloudStorageSinkConnector(w io.Writer, projectID, region, connectClusterID, connectorID, topics, gcsBucketName, tasksMax, formatOutputType, valueConverter, valueConverterSchemasEnable, keyConverter, gcsCredentialsDefault string, opts ...option.ClientOption) error {
	// TODO(developer): Update with your config values. Here is a sample configuration:
	// projectID := "my-project-id"
	// region := "us-central1"
	// connectClusterID := "my-connect-cluster"
	// connectorID := "GCS_SINK_CONNECTOR_ID"
	// topics := "GMK_TOPIC_ID"
	// gcsBucketName := "GCS_BUCKET_NAME"
	// tasksMax := "3"
	// formatOutputType := "json"
	// valueConverter := "org.apache.kafka.connect.json.JsonConverter"
	// valueConverterSchemasEnable := "false"
	// keyConverter := "org.apache.kafka.connect.storage.StringConverter"
	// gcsCredentialsDefault := "true"
	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)

	config := map[string]string{
		"connector.class":                "io.aiven.kafka.connect.gcs.GcsSinkConnector",
		"tasks.max":                      tasksMax,
		"topics":                         topics,
		"gcs.bucket.name":                gcsBucketName,
		"gcs.credentials.default":        gcsCredentialsDefault,
		"format.output.type":             formatOutputType,
		"name":                           connectorID,
		"value.converter":                valueConverter,
		"value.converter.schemas.enable": valueConverterSchemasEnable,
		"key.converter":                  keyConverter,
	}

	connector := &managedkafkapb.Connector{
		Name:    fmt.Sprintf("%s/connectors/%s", parent, connectorID),
		Configs: config,
	}

	req := &managedkafkapb.CreateConnectorRequest{
		Parent:      parent,
		ConnectorId: connectorID,
		Connector:   connector,
	}

	resp, err := client.CreateConnector(ctx, req)
	if err != nil {
		return fmt.Errorf("client.CreateConnector got err: %w", err)
	}
	fmt.Fprintf(w, "Created Cloud Storage sink connector: %s\n", resp.Name)
	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.ConnectorName;
import com.google.cloud.managedkafka.v1.CreateConnectorRequest;
import com.google.cloud.managedkafka.v1.ManagedKafkaConnectClient;
import java.io.IOException;
import java.util.HashMap;
import java.util.Map;

public class CreateCloudStorageSinkConnector {

  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 connectClusterId = "my-connect-cluster";
    String connectorId = "my-gcs-sink-connector";
    String bucketName = "my-gcs-bucket";
    String kafkaTopicName = "kafka-topic";
    String connectorClass = "io.aiven.kafka.connect.gcs.GcsSinkConnector";
    String maxTasks = "3";
    String gcsCredentialsDefault = "true";
    String formatOutputType = "json";
    String valueConverter = "org.apache.kafka.connect.json.JsonConverter";
    String valueSchemasEnable = "false";
    String keyConverter = "org.apache.kafka.connect.storage.StringConverter";
    createCloudStorageSinkConnector(
        projectId,
        region,
        connectClusterId,
        connectorId,
        bucketName,
        kafkaTopicName,
        connectorClass,
        maxTasks,
        gcsCredentialsDefault,
        formatOutputType,
        valueConverter,
        valueSchemasEnable,
        keyConverter);
  }

  public static void createCloudStorageSinkConnector(
      String projectId,
      String region,
      String connectClusterId,
      String connectorId,
      String bucketName,
      String kafkaTopicName,
      String connectorClass,
      String maxTasks,
      String gcsCredentialsDefault,
      String formatOutputType,
      String valueConverter,
      String valueSchemasEnable,
      String keyConverter)
      throws Exception {

    // Build the connector configuration
    Map<String, String> configMap = new HashMap<>();
    configMap.put("connector.class", connectorClass);
    configMap.put("tasks.max", maxTasks);
    configMap.put("topics", kafkaTopicName);
    configMap.put("gcs.bucket.name", bucketName);
    configMap.put("gcs.credentials.default", gcsCredentialsDefault);
    configMap.put("format.output.type", formatOutputType);
    configMap.put("name", connectorId);
    configMap.put("value.converter", valueConverter);
    configMap.put("value.converter.schemas.enable", valueSchemasEnable);
    configMap.put("key.converter", keyConverter);

    Connector connector = Connector.newBuilder()
        .setName(
            ConnectorName.of(projectId, region, connectClusterId, connectorId).toString())
        .putAllConfigs(configMap)
        .build();

    try (ManagedKafkaConnectClient managedKafkaConnectClient = ManagedKafkaConnectClient.create()) {
      CreateConnectorRequest request = CreateConnectorRequest.newBuilder()
          .setParent(ConnectClusterName.of(projectId, region, connectClusterId).toString())
          .setConnectorId(connectorId)
          .setConnector(connector)
          .build();

      // This operation is being handled synchronously.
      Connector response = managedKafkaConnectClient.createConnector(request);
      System.out.printf("Created Cloud Storage Sink connector: %s\n", response.getName());
    } catch (IOException | ApiException e) {
      System.err.printf("managedKafkaConnectClient.createConnector 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.managedkafka_v1.types import Connector, CreateConnectorRequest

connect_client = ManagedKafkaConnectClient()
parent = connect_client.connect_cluster_path(project_id, region, connect_cluster_id)

configs = {
    "connector.class": "io.aiven.kafka.connect.gcs.GcsSinkConnector",
    "tasks.max": tasks_max,
    "topics": topics,
    "gcs.bucket.name": gcs_bucket_name,
    "gcs.credentials.default": "true",
    "format.output.type": format_output_type,
    "name": connector_id,
    "value.converter": value_converter,
    "value.converter.schemas.enable": value_converter_schemas_enable,
    "key.converter": key_converter,
}

connector = Connector()
connector.name = connector_id
connector.configs = configs

request = CreateConnectorRequest(
    parent=parent,
    connector_id=connector_id,
    connector=connector,
)

try:
    operation = connect_client.create_connector(request=request)
    print(f"Waiting for operation {operation.operation.name} to complete...")
    response = operation.result()
    print("Created Connector:", response)
except GoogleAPICallError as e:
    print(f"The operation failed with error: {e}")

创建连接器后,您可以修改、删除、暂停、停止或重启连接器。

接下来怎么做?

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