Cloud Storage Sink コネクタを作成する

Cloud Storage Sink コネクタは、Kafka から Cloud Storage バケットにデータをストリーミングします。これにより、大量のデータを費用対効果の高いスケーラブルな方法で保存して処理できます。

Cloud Storage Sink コネクタのユースケースは次のとおりです。

  • データレイクの取り込み。長期保存とバッチ処理のために、Kafka データをデータレイクに保存します。

  • 規制要件を満たすためにデータをアーカイブします。

始める前に

Cloud Storage Sink コネクタを作成する前に、次のものがあることを確認してください。

必要なロールと権限

コネクタの作成に必要な権限を取得するには、プロジェクトに対するManaged Kafka コネクタ編集者 roles/managedkafka.connectorEditor)IAM ロールの付与を管理者に依頼してください。ロールの付与については、プロジェクト、フォルダ、組織に対するアクセス権の管理をご覧ください。

この事前定義ロールには コネクタの作成に必要な権限が含まれています。必要とされる正確な権限については、「必要な権限」セクションを開いてご確認ください。

必要な権限

コネクタを作成するには、次の権限が必要です。

  • コネクタを作成する: managedkafka.connectors.create

カスタムロールや他の事前定義ロールを使用して、これらの権限を取得することもできます。

Cloud Storage バケットに書き込む権限を付与する

Managed Kafka サービス アカウントには、Cloud Storage バケットにメッセージを書き込むための次の権限が必要です。

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

Cloud Storage バケットを含むプロジェクトのサービス アカウントにストレージ オブジェクト ユーザー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 Sink コネクタの仕組み

Cloud Storage Sink コネクタは、1 つ以上の Kafka トピックからデータを pull し、そのデータを単一の Cloud Storage バケット内のオブジェクトに書き込みます。

Cloud Storage Sink コネクタがデータをコピーする仕組みの詳細を次に示します。

  • コネクタは、ソース クラスタ内の 1 つ以上の Kafka トピックからメッセージを消費します。

  • コネクタは、コネクタ構成で指定したターゲット Cloud Storage バケットにデータを書き込みます。

  • コネクタは、コネクタ構成の特定のプロパティを参照して、Cloud Storage バケットに書き込むときにデータの形式を設定します。 デフォルトでは、出力ファイルは CSV 形式です。format.output.type プロパティを構成して、JSON などのさまざまな出力形式を指定できます。

  • コネクタは、Cloud Storage バケットに書き込まれるファイルの名前も付けます。file.name.prefix プロパティと file.name.template プロパティを使用して、ファイル名をカスタマイズできます。たとえば、ファイル名に Kafka トピック名やメッセージキーを含めることができます。

  • Kafka レコードには、ヘッダー、キー、値の 3 つのコンポーネントがあります。

    • ヘッダーを含めるように format.output.fields を設定すると、出力ファイルにヘッダーを含めることができます。 例: format.output.fields=value,headers

    • key を含めるように format.output.fields を設定すると、出力ファイルにキーを含めることができます。例: 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.typecsv に設定されている場合、キー コンバータは org.apache.kafka.connect.converters.ByteArrayConverter または org.apache.kafka.connect.storage.StringConverter である必要があり、値コンバータは org.apache.kafka.connect.converters.ByteArrayConverter である必要があります。

    • format.output.typejson に設定されている場合、value.converter.schemas.enable プロパティが true であっても、値とキーのスキーマは出力ファイル内のデータとともに書き込まれません。

  • tasks.max プロパティは、コネクタの並列処理レベルを制御します。tasks.max を増やすとスループットが向上しますが、実際の並列処理は Kafka トピックのパーティション数によって制限されます。

Cloud Storage Sink コネクタのプロパティ

Cloud Storage Sink コネクタを作成するときは、次のプロパティを指定します。

コネクタ名

コネクタの名前または ID。リソースの命名方法のガイドラインについては、 Managed Service for Apache Kafka リソースの命名ガイドラインをご覧ください。 名前は変更できません。

コネクタ プラグインのタイプ

コネクタ プラグインのタイプとして [Cloud Storage Sink] を Google Cloud コンソールで選択します。ユーザー インターフェースを使用してコネクタを構成しない場合は、コネクタ クラスも指定する必要があります。

トピック

コネクタがメッセージを消費する Kafka トピック。 1 つ以上のトピックを指定することも、正規表現を使用して複数のトピックを照合することもできます。たとえば、topic.* は「topic」で始まるすべてのトピックに一致します。これらのトピックは、Connect クラスタに関連付けられた Managed Service for Apache Kafka クラスタ内に存在する必要があります。

Cloud Storage バケット

データを保存する Cloud Storage バケットを選択または作成します。

構成

このセクションでは、Cloud Storage Sink コネクタのコネクタ固有の追加の構成プロパティを指定できます。

Kafka トピックのデータは、Avro、JSON、未加工のバイトなど、さまざまな形式で保存できるため、構成の重要な部分としてコンバータの指定があります。 コンバータは、Kafka トピックで使用される形式から Kafka Connect の標準化された内部形式にデータを変換します。Cloud Storage Sink コネクタは、この内部データを受け取り、書き込む前に Cloud Storage バケットに必要な形式に変換します。

Kafka Connect でのコンバータの役割、 サポートされているコンバータのタイプ、一般的な構成オプションの詳細については、 コンバータをご覧ください。

Cloud Storage Sink コネクタに固有の構成を次に示します。

  • 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 Sink コネクタの構成をご覧ください。

Cloud Storage Sink コネクタを作成する

コネクタを作成する前に、Cloud Storage Sink コネクタのプロパティのドキュメントを確認してください。

コンソール

  1. コンソールで、[Connect クラスタ] ページに移動します。 Google Cloud

    [Connect クラスタ] に移動

  2. コネクタを作成する Connect クラスタをクリックします。

    [Connect クラスタの詳細] ページが表示されます。

  3. [コネクタを作成] をクリックします。

    [Kafka コネクタを作成] ページが表示されます。

  4. コネクタ名に文字列を入力します。

    コネクタの命名方法のガイドラインについては、 Managed Service for Apache Kafka リソースの命名ガイドラインをご覧ください。

  5. [**コネクタ プラグイン**] で [**Cloud Storage Sink**] を選択します。

  6. データをストリーミングできるトピック を指定します。

  7. データを保存するストレージ バケット を選択します。

  8. (省略可)[構成] セクションで追加の設定を行います。

  9. タスクの再起動ポリシー を選択します。詳細については、 タスクの再起動ポリシーをご覧ください。

  10. [作成] をクリックします。

gcloud

  1. コンソールで Cloud Shell をアクティブにします。 Google Cloud

    Cloud Shell をアクティブにする

    コンソールの下部にある Google Cloud Cloud Shell セッションが開始し、コマンドライン プロンプトが表示されます。Cloud Shell はシェル環境です 。Google Cloud CLI がすでにインストールされており、現在のプロジェクトの値もすでに設定されています 。セッションが初期化されるまで数秒かかることがあります。

  2. gcloud managed-kafka connectors create コマンドを実行します。

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

    次のように置き換えます。

    Cloud Storage Sink コネクタの構成ファイルの例を次に示します。

    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: データが Cloud Storage Sink コネクタに流れる Managed Service for Apache Kafka トピックの ID。

    • GCS_BUCKET_NAME: パイプラインのシンクとして機能する Cloud Storage バケットの名前。

    • GCS_SINK_CONNECTOR_ID: Cloud Storage Sink コネクタの 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 または米国その他の諸国における関連会社の商標です。