建立 BigQuery Sink 連接器

BigQuery 接收器連接器會將 Kafka 的資料串流至 BigQuery 資料表,方便您在 BigQuery 中擷取及分析即時資料。

BigQuery Sink 連接器的用途包括:

  • 資料倉儲。將串流資料載入 BigQuery,用於分析和報表。

  • 填入即時資訊主頁所用的 BigQuery 資料表。

事前準備

建立 BigQuery Sink 連接器前,請確認您已備妥下列項目:

必要角色和權限

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

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

所需權限

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

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

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

授予寫入 BigQuery 資料表的權限

代管 Kafka 服務帳戶必須具備將訊息寫入 BigQuery 資料表的權限。在包含資料表的專案中,將「BigQuery 資料編輯者」(roles/bigquery.dataEditor) 角色授予服務帳戶。

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

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

重要設定參數

以下是 BigQuery 接收器連接器的設定建議。

佇列大小

queueSize 設定會控管寫入 BigQuery 的要求數量上限。根據預設,系統不會設定這個值,這可能會導致 Connect 叢集發生記憶體不足錯誤。建議您設定這個值,以減少發生這類錯誤的機率。佇列中的項目是寫入要求,不會直接對應至訊息或位元組數量。因此,我們很難事先得知建議值。

從 Kafka 輪詢訊息時,連接器會根據訊息的目的地資料表分割工作。如果投票活動中的所有訊息都要寫入單一表格,則只會新增一個工作項目。但如果輪詢中的訊息要寫入三個資料表,系統就會新增三個工作項目。您可以根據訊息輪詢設定、Kafka 訊息大小,以及要寫入的資料表數量,大致估算工作項目的大小。

BigQuery Sink 連接器的結構定義

BigQuery 接收器連接器會使用設定的值轉換器 (value.converter),將 Kafka 記錄值剖析為欄位。然後將欄位寫入 BigQuery 資料表中同名的資料欄。

連接器需要結構定義才能運作。您可以透過下列方式提供結構定義:

  • 以訊息為準的結構定義:結構定義會納入每則訊息。
  • 以資料表為準的結構定義:連接器會根據 BigQuery 資料表結構定義推論訊息結構定義。
  • 結構定義儲存庫:連接器會從結構定義儲存庫 (例如 Managed Service for Apache Kafka 結構定義儲存庫 (預先發布版)) 讀取結構定義。

以下各節將說明這些選項。

以訊息為基礎的結構定義

在這個模式中,每筆 Kafka 記錄都包含 JSON 結構定義。連接器會使用結構定義,將記錄資料寫入為 BigQuery 資料表列。

如要使用以訊息為準的結構定義,請在連接器上設定下列屬性:

  • value.converter=org.apache.kafka.connect.json.JsonConverter
  • value.converter.schemas.enable=true

Kafka 記錄值範例:

{
  "schema": {
    "type": "struct",
    "fields": [
      {
        "field": "user",
        "type": "string",
        "optional": false
      },
      {
        "field": "age",
        "type": "int64",
        "optional": false
      }
    ]
  },
  "payload": {
    "user": "userId",
    "age": 30
  }
}

如果目的地資料表已存在,BigQuery 資料表結構定義必須與內嵌訊息結構定義相容。如果autoCreateTables=true,連接器會視需要自動建立目的地資料表。詳情請參閱「建立資料表」。

如要讓連接器在訊息結構定義變更時更新 BigQuery 資料表結構定義,請將 allowNewBigQueryFieldsallowSchemaUnionizationallowBigQueryRequiredFieldRelaxation 設為 true

以資料表為基礎的結構定義

在此模式下,Kafka 記錄包含純 JSON 資料,沒有明確的結構定義。連接器會從目的地資料表推論結構定義。

需求條件:

  • BigQuery 資料表必須已存在。
  • Kafka 記錄資料必須與資料表結構定義相容。
  • 這個模式不支援根據傳入訊息動態更新結構定義。

如要使用以表格為基礎的結構定義,請在連接器上設定下列屬性:

  • value.converter=org.apache.kafka.connect.json.JsonConverter
  • value.converter.schemas.enable=false
  • bigQueryPartitionDecorator=false

如果 BigQuery 資料表使用時間分區,且分區頻率為每日,則 bigQueryPartitionDecorator 可以是 true。否則,請將這個屬性設為 false

Kafka 記錄值範例:

{
  "user": "userId",
  "age": 30
}

結構定義儲存庫

在此模式下,每筆 Kafka 記錄都包含 Apache Avro 資料,且訊息結構定義會儲存在結構定義儲存庫中。

如要搭配結構定義登錄使用 BigQuery Sink 連接器,請在連接器上設定下列屬性:

  • value.converter=io.confluent.connect.avro.AvroConverter
  • value.converter.schema.registry.url=SCHEMA_REGISTRY_URL

SCHEMA_REGISTRY_URL 替換為結構定義登錄的網址。

如要搭配 Managed Service for Apache Kafka 結構定義登錄使用連接器,請設定下列屬性:

  • value.converter.bearer.auth.credentials.source=GCP

詳情請參閱「Use Kafka Connect with schema registry」。

Apache Iceberg 代管資料表

BigQuery Sink 連接器支援 Apache Iceberg 代管資料表 (以下簡稱「Iceberg 代管資料表」) 做為接收器目標。

Iceberg 代管資料表是 Google Cloud上開放格式 lakehouse 的建構基礎。Iceberg 代管資料表提供與 BigQuery 資料表相同的全代管體驗,但會使用 Parquet 將資料儲存在客戶擁有的儲存空間 bucket 中,以便與 Apache Iceberg 開放資料表格式互通。

如要瞭解如何建立 Apache Iceberg 資料表,請參閱「建立 Apache Iceberg 資料表」。

建立 BigQuery Sink 連接器

控制台

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

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

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

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

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

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

  5. 針對「連接器外掛程式」,選取「BigQuery 接收器」

  6. 在「主題」部分,指定要讀取的 Kafka 主題。您可以指定主題清單或規則運算式,比對主題名稱。

    • 方法 1:選擇「選取 Kafka 主題清單」。在「Kafka topics」(Kafka 主題) 清單中,選取一或多個主題。然後點選「OK」

    • 方法 2:選擇「使用主題規則運算式」。在「主題規則運算式」欄位中,輸入規則運算式。

  7. 按一下「資料集」,然後指定 BigQuery 資料集。您可以選擇現有資料集或建立新資料集。

  8. 選用:在「設定」方塊中,新增設定屬性或編輯預設屬性。詳情請參閱「設定連接器」。

  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:BigQuery Sink 連接器的 YAML 設定檔路徑。

    以下是 BigQuery Sink 連接器的設定檔範例:

    name: "BQ_SINK_CONNECTOR_ID"
    project: "GCP_PROJECT_ID"
    topics: "GMK_TOPIC_ID"
    tasks.max: 3
    connector.class: "com.wepay.kafka.connect.bigquery.BigQuerySinkConnector"
    key.converter: "org.apache.kafka.connect.storage.StringConverter"
    value.converter: "org.apache.kafka.connect.json.JsonConverter"
    value.converter.schemas.enable: "false"
    defaultDataset: "BQ_DATASET_ID"
    

    更改下列內容:

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

    • GCP_PROJECT_ID:BigQuery 資料集所在的 Google Cloud專案 ID。

    • GMK_TOPIC_ID:資料從中流向 BigQuery 接收器連接器的 Managed Service for Apache Kafka 主題 ID。

    • BQ_DATASET_ID:做為管道接收器的 BigQuery 資料集 ID。

Terraform

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

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

  configs = {
    "name"                           = "my-bigquery-sink-connector"
    "project"                        = data.google_project.default.project_id
    "topics"                         = "GMK_TOPIC_ID"
    "tasks.max"                      = "3"
    "connector.class"                = "com.wepay.kafka.connect.bigquery.BigQuerySinkConnector"
    "key.converter"                  = "org.apache.kafka.connect.storage.StringConverter"
    "value.converter"                = "org.apache.kafka.connect.json.JsonConverter"
    "value.converter.schemas.enable" = "false"
    "defaultDataset"                 = "BQ_DATASET_ID"
  }

  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"
)

// createBigQuerySinkConnector creates a BigQuery Sink connector.
func createBigQuerySinkConnector(w io.Writer, projectID, region, connectClusterID, connectorID, topics, tasksMax, keyConverter, valueConverter, valueConverterSchemasEnable, defaultDataset 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 := "BQ_SINK_CONNECTOR_ID"
	// topics := "GMK_TOPIC_ID"
	// tasksMax := "3"
	// keyConverter := "org.apache.kafka.connect.storage.StringConverter"
	// valueConverter := "org.apache.kafka.connect.json.JsonConverter"
	// valueConverterSchemasEnable := "false"
	// defaultDataset := "BQ_DATASET_ID"
	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)

	// BigQuery Sink sample connector configuration
	config := map[string]string{
		"name":                           connectorID,
		"project":                        projectID,
		"topics":                         topics,
		"tasks.max":                      tasksMax,
		"connector.class":                "com.wepay.kafka.connect.bigquery.BigQuerySinkConnector",
		"key.converter":                  keyConverter,
		"value.converter":                valueConverter,
		"value.converter.schemas.enable": valueConverterSchemasEnable,
		"defaultDataset":                 defaultDataset,
	}

	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 BigQuery 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 CreateBigQuerySinkConnector {

  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-bigquery-sink-connector";
    String bigqueryProjectId = "my-bigquery-project-id";
    String datasetName = "my-dataset";
    String kafkaTopicName = "kafka-topic";
    String maxTasks = "3";
    String connectorClass = "com.wepay.kafka.connect.bigquery.BigQuerySinkConnector";
    String keyConverter = "org.apache.kafka.connect.storage.StringConverter";
    String valueConverter = "org.apache.kafka.connect.json.JsonConverter";
    String valueSchemasEnable = "false";
    createBigQuerySinkConnector(
        projectId,
        region,
        connectClusterId,
        connectorId,
        bigqueryProjectId,
        datasetName,
        kafkaTopicName,
        maxTasks,
        connectorClass,
        keyConverter,
        valueConverter,
        valueSchemasEnable);
  }

  public static void createBigQuerySinkConnector(
      String projectId,
      String region,
      String connectClusterId,
      String connectorId,
      String bigqueryProjectId,
      String datasetName,
      String kafkaTopicName,
      String maxTasks,
      String connectorClass,
      String keyConverter,
      String valueConverter,
      String valueSchemasEnable)
      throws Exception {

    // Build the connector configuration
    Map<String, String> configMap = new HashMap<>();
    configMap.put("name", connectorId);
    configMap.put("project", bigqueryProjectId);
    configMap.put("topics", kafkaTopicName);
    configMap.put("tasks.max", maxTasks);
    configMap.put("connector.class", connectorClass);
    configMap.put("key.converter", keyConverter);
    configMap.put("value.converter", valueConverter);
    configMap.put("value.converter.schemas.enable", valueSchemasEnable);
    configMap.put("defaultDataset", datasetName);

    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 BigQuery 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 = {
    "name": connector_id,
    "project": project_id,
    "topics": topics,
    "tasks.max": tasks_max,
    "connector.class": "com.wepay.kafka.connect.bigquery.BigQuerySinkConnector",
    "key.converter": key_converter,
    "value.converter": value_converter,
    "value.converter.schemas.enable": value_converter_schemas_enable,
    "defaultDataset": default_dataset,
}

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}")

建立連接器後,您可以編輯、刪除、暫停、停止或重新啟動連接器。

設定連接器

本節說明您可以在連接器上設定的部分設定屬性。如需這個連接器專屬屬性的完整清單,請參閱 BigQuery Sink 連接器設定

資料表名稱

根據預設,連接器會使用主題名稱做為 BigQuery 資料表名稱。如要使用其他資料表名稱,請設定 topic2TableMap 屬性,格式如下:

topic2TableMap=TOPIC_1:TABLE_1,TOPIC_2:TABLE_2,...

建立資料表

如果目的地資料表不存在,BigQuery Sink 連接器可以建立這些資料表。

  • 如果 autoCreateTables=true,連接器會嘗試建立不存在的 BigQuery 資料表。這是預設行為。

  • 如果為 autoCreateTables=false,連接器不會建立任何資料表。如果目的地資料表不存在,就會發生錯誤。

如果 autoCreateTablestrue,您可以使用下列設定屬性,更精細地控管連接器建立及設定新資料表的方式:

  • allBQFieldsNullable
  • clusteringPartitionFieldNames
  • convertDoubleSpecialValues
  • partitionExpirationMs
  • sanitizeFieldNames
  • sanitizeTopics
  • timestampPartitionFieldName

如要瞭解這些屬性,請參閱「BigQuery Sink 連接器設定」。

Kafka 中繼資料

您可以分別設定 kafkaDataFieldNamekafkaKeyFieldName 欄位,將 Kafka 中的其他資料 (例如中繼資料資訊和鍵值資訊) 對應至 BigQuery 資料表。中繼資料資訊的例子包括 Kafka 主題、分割區、偏移和插入時間。

後續步驟

Apache Kafka® 是 The Apache Software Foundation 或其關聯企業在美國與/或其他國家/地區的註冊商標。