建立 Pub/Sub Sink 連接器

Pub/Sub 接收器連接器會將 Kafka 主題的訊息串流至 Pub/Sub 主題,方便您將以 Kafka 為基礎的應用程式與 Pub/Sub 整合。

Pub/Sub 接收器連接器的用途包括:

  • 事件導向架構。從發布至 Kafka 的訊息觸發 Pub/Sub 處理作業。

  • 服務整合。將 Kafka 的資料傳送至其他 Google Cloud 服務 或應用程式,這些服務或應用程式會從 Pub/Sub 取用資料。

  • 根據處理的資料觸發即時通知或動作。

事前準備

建立 Pub/Sub 接收器連接器前,請確認您已備妥下列項目:

必要角色和權限

如要取得建立連線器所需的權限,請要求管理員授予您專案的「Managed Kafka Connector Editor」(代管 Kafka 連線器編輯者) (roles/managedkafka.connectorEditor) IAM 角色。如要進一步瞭解如何授予角色,請參閱「管理專案、資料夾和組織的存取權」。

這個預先定義的角色具備建立連接器所需的權限。如要查看確切的必要權限,請展開「Required permissions」(必要權限) 部分:

所需權限

如要建立連接器,必須具備下列權限:

  • 建立連接器: managedkafka.connectors.create

您或許還可透過自訂角色或其他預先定義的角色取得這些權限。

授予發布至 Pub/Sub 主題的權限

代管 Kafka 服務帳戶必須具備將訊息發布至 Pub/Sub 主題的權限。在包含 Pub/Sub 主題的專案中,將 Pub/Sub 發布者 (roles/pubsub.publisher) 角色授予服務帳戶。

代管 Kafka 服務帳戶的格式如下: service-PROJECT_NUMBER@gcp-sa-managedkafka.iam.gserviceaccount.com, 其中 PROJECT_NUMBER 是 Connect 叢集的專案編號。

如果 Connect 叢集與 Managed Service for Apache Kafka 叢集位於不同專案,請參閱「 在不同專案中建立 Connect 叢集」。

Pub/Sub 接收器連接器的運作方式

Pub/Sub 接收器連接器會從一或多個 Kafka 主題提取訊息,並發布至 Pub/Sub 主題。

以下詳細說明 Pub/Sub Sink 連接器如何複製資料:

  • 連接器會取用來源叢集內一或多個 Kafka 主題的訊息。

  • 連結器會將訊息寫入 cps.topic 設定屬性指定的目標 Pub/Sub 主題 ID。這是必填屬性。

  • 連結器也需要使用 cps.project 設定屬性,指定包含 Pub/Sub 主題的 Google Cloud 專案。這是必要屬性。

  • 連接器也可以選擇使用自訂 Pub/Sub 端點,方法是使用 cps.endpoint 屬性指定端點。預設端點為 "pubsub.googleapis.com:443"

  • 為提升效能,連接器會先緩衝處理訊息,再發布至 Pub/Sub。您可以設定 maxBufferSizemaxBufferBytesmaxDelayThresholdMsmaxOutstandingRequestBytesmaxOutstandingMessages,控制緩衝。

  • Kafka 記錄包含三個元件:標頭、鍵和值。 連接器會使用鍵值轉換器,將 Kafka 訊息資料轉換為 Pub/Sub 預期的格式。使用結構體或對應值結構定義時,messageBodyName 屬性會指定要用做 Pub/Sub 訊息內文的欄位或鍵。

  • 只要將 metadata.publish 屬性設為 true,連接器就能將 Kafka 主題、分割區、位移和時間戳記做為訊息屬性。

  • 連接器可將 Kafka 訊息標頭納入 Pub/Sub 訊息屬性,方法是將 headers.publish 屬性設為 true

  • 連接器可以使用 orderingKeySource 屬性,為 Pub/Sub 訊息加入排序鍵。可能的值包括 "none" (預設值)、"key""partition"

  • tasks.max 屬性可控制連接器的平行處理層級。增加 tasks.max 可提高輸送量,但實際的平行處理量會受 Kafka 主題中的分區數量限制。

Pub/Sub 接收器連接器的屬性

建立 Pub/Sub 接收器連接器時,您需要指定下列屬性。

連接器名稱

Connect 叢集內連接器的專屬名稱。 如要查看資源命名準則,請參閱 Managed Service for Apache Kafka 資源命名指南

連接器外掛程式類型

選取「Pub/Sub Sink」做為連接器外掛程式類型。這會決定資料流程的方向 (從 Kafka Pub/Sub),以及使用的特定連接器實作方式。如果您未使用使用者介面設定連接器,則必須一併指定連接器類別。

Kafka 主題

連接器從中取用訊息的 Kafka 主題。 您可以指定一或多個主題,也可以使用規則運算式比對多個主題。例如,topic.* 可比對所有以「topic」開頭的主題。這些主題必須位於與 Connect 叢集相關聯的 Managed Service for Apache Kafka 叢集內。

Pub/Sub 主題

現有的 Pub/Sub 主題,連接器會將訊息發布至該主題。如「事前準備」一文所述,請確保 Connect 叢集服務帳戶具備主題專案的 roles/pubsub.publisher 角色。

設定

您可以在這個部分指定其他連接器專屬的設定屬性。

由於 Kafka 主題中的資料格式可能不一,例如 Avro、JSON 或原始位元組,因此設定的關鍵部分是指定轉換器。轉換器會將 Kafka 主題中使用的格式資料,轉換為 Kafka Connect 的標準內部格式。接著,Pub/Sub Sink 連接器會取得這項內部資料,並轉換成 Pub/Sub 要求的格式,然後寫入資料。

如要進一步瞭解 Kafka Connect 中轉換器的角色、支援的轉換器類型和常見設定選項,請參閱轉換器

以下是 Pub/Sub 接收器連接器的專屬設定:

  • cps.project:指定包含 Pub/Sub 主題的 Google Cloud 專案 ID。

  • cps.topic:指定要將資料發布至哪個 Pub/Sub 主題。

  • cps.endpoint:指定要使用的 Pub/Sub 端點。

如要查看這個連接器可用的特定設定屬性清單,請參閱「Pub/Sub Sink 連接器設定」。

建立 Pub/Sub 接收器連接器

建立連接器前,請先參閱 Pub/Sub 接收器連接器屬性的說明文件。

控制台

  1. 前往 Google Cloud 控制台的「Connect Clusters」(連結叢集) 頁面。

    前往「Connect Clusters」(連結叢集)

  2. 按一下要建立連接器的 Connect 叢集。

    系統隨即會顯示「Connect cluster details」(連結叢集詳細資料) 頁面。

  3. 按一下「Create connector」(建立連接器)。

    系統會顯示「建立 Kafka 連接器」頁面。

  4. 輸入連接器名稱字串。

    如要查看連接器命名準則,請參閱 Managed Service for Apache Kafka 資源命名指南

  5. 在「連接器外掛程式」部分,選取「Pub/Sub 接收器」

  6. 在「主題」下方,選擇「選取 Kafka 主題清單」或「使用主題規則運算式」。然後選取或輸入這個連接器要取用訊息的 Kafka 主題。這些主題位於相關聯的 Kafka 叢集中。

  7. 在「Select a Cloud Pub/Sub topic」(選取 Cloud Pub/Sub 主題) 中,選擇這個連接器要將訊息發布至哪個 Pub/Sub 主題。主題會以完整資源名稱格式顯示:projects/{project}/topics/{topic}

  8. (選用) 在「設定」部分調整其他設定。您可以在這裡指定 tasks.maxkey.convertervalue.converter 等屬性,如上一節所述。

  9. 選取「任務重新啟動政策」。詳情請參閱「工作重新啟動政策」。

  10. 點選「建立」

gcloud

  1. 在 Google Cloud 控制台中啟用 Cloud Shell。

    啟用 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
    

    更改下列內容:

    • CONNECTOR_ID:連接器的 ID 或名稱。 如要查看連線器命名準則,請參閱 Managed Service for Apache Kafka 資源命名指南。 連接器名稱無法變更。

    • LOCATION:建立連接器的位置。這個位置必須與您建立 Connect 叢集的位置相同。

    • CONNECT_CLUSTER_ID:建立連接器的 Connect 叢集 ID。

    • CONFIG_FILE:連接器的 YAML 設定檔路徑。

    以下是 Pub/Sub 接收器連接器的設定檔範例:

    connector.class: "com.google.pubsub.kafka.sink.CloudPubSubSinkConnector"
    name: "CPS_SINK_CONNECTOR_ID"
    tasks.max: "1"
    topics: "GMK_TOPIC_ID"
    value.converter: "org.apache.kafka.connect.storage.StringConverter"
    key.converter: "org.apache.kafka.connect.storage.StringConverter"
    cps.topic: "CPS_TOPIC_ID"
    cps.project: "GCP_PROJECT_ID"
    

    更改下列內容:

    • CPS_SINK_CONNECTOR_ID:Pub/Sub Sink 連接器的 ID 或名稱。如要查看連線器命名準則,請參閱 Managed Service for Apache Kafka 資源命名指南。連接器名稱無法變更。

    • GMK_TOPIC_ID:Pub/Sub Sink 連接器從中讀取資料的 Managed Service for Apache Kafka 主題 ID。

    • CPS_TOPIC_ID:資料發布至其中的 Pub/Sub 主題 ID。

    • GCP_PROJECT_ID:Pub/Sub 主題所在的專案 ID。 Google Cloud

Terraform

您可以使用 Terraform 資源建立連接器

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

  configs = {
    "connector.class" = "com.google.pubsub.kafka.sink.CloudPubSubSinkConnector"
    "name"            = "my-pubsub-sink-connector"
    "tasks.max"       = "3"
    "topics"          = "TOPIC_NAME"
    "cps.topic"       = "CPS_TOPIC_NAME"
    "cps.project"     = "CPS_PROJECT_NAME"
    "value.converter" = "org.apache.kafka.connect.storage.StringConverter"
    "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"
)

// createPubSubSinkConnector creates a Pub/Sub Sink connector.
func createPubSubSinkConnector(w io.Writer, projectID, region, connectClusterID, connectorID, topics, valueConverter, keyConverter, cpsTopic, cpsProject, tasksMax 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 := "CPS_SINK_CONNECTOR_ID"
	// topics := "GMK_TOPIC_ID"
	// valueConverter := "org.apache.kafka.connect.storage.StringConverter"
	// keyConverter := "org.apache.kafka.connect.storage.StringConverter"
	// cpsTopic := "CPS_TOPIC_ID"
	// cpsProject := "GCP_PROJECT_ID"
	// tasksMax := "3"
	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)

	// Pub/Sub Sink sample connector configuration
	config := map[string]string{
		"connector.class": "com.google.pubsub.kafka.sink.CloudPubSubSinkConnector",
		"name":            connectorID,
		"tasks.max":       tasksMax,
		"topics":          topics,
		"value.converter": valueConverter,
		"key.converter":   keyConverter,
		"cps.topic":       cpsTopic,
		"cps.project":     cpsProject,
	}

	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 Pub/Sub 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 CreatePubSubSinkConnector {

  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-pubsub-sink-connector";
    String pubsubProjectId = "my-pubsub-project-id";
    String pubsubTopicName = "my-pubsub-topic";
    String kafkaTopicName = "kafka-topic";
    String connectorClass = "com.google.pubsub.kafka.sink.CloudPubSubSinkConnector";
    String maxTasks = "3";
    String valueConverter = "org.apache.kafka.connect.storage.StringConverter";
    String keyConverter = "org.apache.kafka.connect.storage.StringConverter";
    createPubSubSinkConnector(
        projectId,
        region,
        connectClusterId,
        connectorId,
        pubsubProjectId,
        pubsubTopicName,
        kafkaTopicName,
        connectorClass,
        maxTasks,
        valueConverter,
        keyConverter);
  }

  public static void createPubSubSinkConnector(
      String projectId,
      String region,
      String connectClusterId,
      String connectorId,
      String pubsubProjectId,
      String pubsubTopicName,
      String kafkaTopicName,
      String connectorClass,
      String maxTasks,
      String valueConverter,
      String keyConverter)
      throws Exception {

    // Build the connector configuration
    Map<String, String> configMap = new HashMap<>();
    configMap.put("connector.class", connectorClass);
    configMap.put("name", connectorId);
    configMap.put("tasks.max", maxTasks);
    configMap.put("topics", kafkaTopicName);
    configMap.put("value.converter", valueConverter);
    configMap.put("key.converter", keyConverter);
    configMap.put("cps.topic", pubsubTopicName);
    configMap.put("cps.project", pubsubProjectId);

    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 Pub/Sub 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": "com.google.pubsub.kafka.sink.CloudPubSubSinkConnector",
    "name": connector_id,
    "tasks.max": tasks_max,
    "topics": topics,
    "value.converter": value_converter,
    "key.converter": key_converter,
    "cps.topic": cps_topic,
    "cps.project": cps_project,
}

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® 是 The Apache Software Foundation 或其關聯企業在美國與/或其他國家/地區的註冊商標。