搭配結構定義註冊資料庫使用 Kafka Connect

本教學課程說明如何整合 Kafka Connect結構定義儲存庫,讓連接器使用已註冊的結構定義序列化及還原序列化訊息。

開始本教學課程前,請先完成「使用結構定義登錄產生 Avro 訊息」中的步驟。完成後,您會取得本教學課程所需的下列元件:

  • 已註冊 Avro 結構定義的結構定義儲存庫。
  • Kafka 叢集,其中包含 Avro 訊息的主題。

總覽

結構定義會定義訊息或資料酬載的結構。結構定義儲存庫提供集中式位置來儲存結構定義。使用結構定義登錄檔,可確保資料從生產者傳輸至消費者時,編碼和解碼方式一致。

本教學課程將示範下列情境:

  1. Java 生產者用戶端會將 Apache Avro 訊息寫入 Managed Service for Apache Kafka 叢集。Avro 訊息結構定義會儲存在結構定義登錄中。

  2. BigQuery 接收器連接器會讀取 Kafka 的傳入訊息,並寫入 BigQuery。連接器會自動建立與 Avro 結構定義相容的 BigQuery 資料表。

事前準備

  1. 登入 Google Cloud 帳戶。如果您是 Google Cloud新手,歡迎 建立帳戶,親自評估產品在實際工作環境中的成效。新客戶還能獲得價值 $300 美元的免費抵免額,可用於執行、測試及部署工作負載。
  2. In the Google Cloud console, on the project selector page, select or create a Google Cloud project.

    Roles required to select or create a project

    • Select a project: Selecting a project doesn't require a specific IAM role—you can select any project that you've been granted a role on.
    • Create a project: To create a project, you need the Project Creator role (roles/resourcemanager.projectCreator), which contains the resourcemanager.projects.create permission. Learn how to grant roles.

    Go to project selector

  3. Verify that billing is enabled for your Google Cloud project.

  4. Enable the Managed Kafka API.

    Roles required to enable APIs

    To enable APIs, you need the Service Usage Admin IAM role (roles/serviceusage.serviceUsageAdmin), which contains the serviceusage.services.enable permission. Learn how to grant roles.

    Enable the API

  5. In the Google Cloud console, on the project selector page, select or create a Google Cloud project.

    Roles required to select or create a project

    • Select a project: Selecting a project doesn't require a specific IAM role—you can select any project that you've been granted a role on.
    • Create a project: To create a project, you need the Project Creator role (roles/resourcemanager.projectCreator), which contains the resourcemanager.projects.create permission. Learn how to grant roles.

    Go to project selector

  6. Verify that billing is enabled for your Google Cloud project.

  7. Enable the Managed Kafka API.

    Roles required to enable APIs

    To enable APIs, you need the Service Usage Admin IAM role (roles/serviceusage.serviceUsageAdmin), which contains the serviceusage.services.enable permission. Learn how to grant roles.

    Enable the API

  8. 請確認您在專案中具備下列一或多個角色: 代管 Kafka 叢集編輯者代管 Kafka Connect 叢集編輯者代管 Kafka 連接器編輯者、 以及 BigQuery 資料擁有者

    檢查角色

    1. 前往 Google Cloud 控制台的「IAM」頁面。

      前往「IAM」頁面
    2. 選取專案。
    3. 在「主體」欄中,找出所有識別您或您所屬群組的資料列。如要瞭解自己所屬的群組,請與管理員聯絡。

    4. 針對指定或包含您的所有列,請檢查「角色」欄,確認角色清單是否包含必要角色。

    授予角色

    1. 前往 Google Cloud 控制台的「IAM」頁面。

      前往「IAM」頁面
    2. 選取專案。
    3. 按一下「Grant access」(授予存取權)
    4. 在「New principals」(新增主體) 欄位中,輸入您的使用者 ID。 這通常是指 Google 帳戶的電子郵件地址。

    5. 按一下「選取角色」,然後搜尋角色。
    6. 如要授予其他角色,請按一下「Add another role」(新增其他角色),然後新增其他角色。
    7. 按一下「Save」(儲存)
  9. 完成「 使用結構定義登錄檔產生 Avro 訊息」中的步驟。

建立連結叢集

如要建立 Connect 叢集,請按照下列步驟操作。建立 Connect 叢集最多可能需要 30 分鐘。

控制台

  1. 前往「Managed Service for Apache Kafka」>「Connect Clusters」(連線叢集) 頁面。

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

  2. 點選 「Create」(建立)

  3. 在「Connect cluster name」(Connect 叢集名稱) 中輸入字串。例如:my-connect-cluster

  4. 在「主要 Kafka 叢集」中,選取 Kafka 叢集。

  5. 點選「建立」

建立 Connect 叢集時,叢集狀態為 Creating。 叢集建立完成後,狀態會顯示為 Active

gcloud

如要建立 Connect 叢集,請執行 gcloud managed-kafka connect-clusters create 指令。

gcloud managed-kafka connect-clusters create CONNECT_CLUSTER \
  --location=REGION \
  --cpu=12 \
  --memory=12GiB \
  --primary-subnet=projects/PROJECT_ID/regions/REGION/subnetworks/SUBNET_NAME \
  --kafka-cluster=KAFKA_CLUSTER \
  --async

更改下列內容:

  • CONNECT_CLUSTER:Connect 叢集名稱
  • REGION:建立 Kafka 叢集的區域
  • PROJECT_ID:專案 ID
  • SUBNET_NAME:建立 Kafka 叢集的子網路
  • KAFKA_CLUSTER:Kafka 叢集名稱

這項指令會以非同步方式執行,並傳回作業 ID:

Check operation [projects/PROJECT_ID/locations/REGION/operations/OPERATION_ID] for status.

如要追蹤建立作業的進度,請使用 gcloud managed-kafka operations describe 指令:

gcloud managed-kafka operations describe OPERATION_ID \
  --location=REGION

詳情請參閱「監控叢集建立作業」。

授予 IAM 角色

將「BigQuery 資料編輯者」身分與存取權管理 (IAM) 角色授予代管 Kafka 服務代理程式。這個角色可讓連接器寫入 BigQuery。

控制台

  1. 前往 Google Cloud 控制台的「IAM」(身分與存取權管理) 頁面。

    前往「身分與存取權管理」頁面

  2. 選取「包含 Google 提供的角色授予項目」

  3. 找到「Managed Kafka Service Account」(受管理 Kafka 服務帳戶) 列,然後按一下 「Edit principal」(編輯主體)

  4. 點選「新增其他角色」,然後選取「BigQuery 資料編輯者」角色。

  5. 按一下 [儲存]

如要進一步瞭解如何授予角色,請參閱「使用控制台授予 IAM 角色」。

gcloud

如要將 IAM 角色授予服務帳戶,請執行 gcloud projects add-iam-policy-binding 指令。

gcloud projects add-iam-policy-binding PROJECT_ID \
    --member=serviceAccount:service-PROJECT_NUMBER@gcp-sa-managedkafka.iam.gserviceaccount.com \
    --role=roles/bigquery.dataEditor

更改下列內容:

  • PROJECT_ID:專案 ID
  • PROJECT_NUMBER:您的專案編號

如要找出專案編號,請使用 gcloud projects describe 指令。

根據預設,Connect 叢集有權存取同一專案中的結構定義登錄。如果您在與結構定義儲存庫不同的專案中建立 Connect 叢集,則必須將結構定義儲存庫專案的「Managed Kafka Client」(roles/managedkafka.client) 角色,授予 Managed Kafka 服務代理。詳情請參閱「在不同專案中建立 Connect 叢集」。

建立 BigQuery 資料集

在這個步驟中,您會建立用來保存 BigQuery 資料表的資料集。BigQuery Sink 連接器會自動建立資料表。

如要建立資料集,請按照下列步驟操作。

控制台

  1. 開啟「BigQuery」BigQuery頁面。

    前往 BigQuery 頁面

  2. 在「Explorer」面板中,選取要建立資料集的專案。

  3. 展開 「查看動作」選項,然後點選「建立資料集」

  4. 在「Create dataset」(建立資料集) 頁面:

    • 在「Dataset ID」(資料集 ID) 中輸入資料集名稱。

    • 針對「位置類型」,選擇資料集的地理位置。

gcloud

如要建立新的資料集,請使用 bq mk 指令,並加上 --dataset 旗標。

bq mk --location REGION \
  --dataset PROJECT_ID:DATASET_NAME

更改下列內容:

  • PROJECT_ID:專案 ID
  • DATASET_NAME:資料集名稱
  • REGION:資料集的位置

建立 BigQuery Sink 連接器

在這個步驟中,您會建立 BigQuery 接收器連接器,從 Kafka 讀取資料並寫入 BigQuery。您可設定連接器,使用儲存在結構定義登錄中的 Avro 結構定義,還原 Kafka 訊息的序列化狀態。

控制台

  1. 前往「Managed Service for Apache Kafka」>「Connect Clusters」(連線叢集) 頁面。

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

  2. 按一下 Connect 叢集名稱。

  3. 按一下 「建立連接器」

  4. 在「Connector name」(連接器名稱) 中輸入字串。範例:bigquery-connector

  5. 在「連接器外掛程式」清單中,選取「BigQuery Sink」。

  6. 在「Topics」(主題) 中,選取名為 newUsers 的 Kafka 主題。這個主題是由 Java 生產端用戶端建立。

  7. 在「資料集」中,輸入 BigQuery 資料集的名稱, 格式如下: PROJECT_ID.DATASET_NAME。 範例:my-project.dataset1

  8. 在「Configurations」(設定) 編輯方塊中,將現有設定替換為下列設定:

    tasks.max=3
    key.converter=org.apache.kafka.connect.storage.StringConverter
    value.converter=io.confluent.connect.avro.AvroConverter
    value.converter.bearer.auth.credentials.source=GCP
    value.converter.schema.registry.url=https://managedkafka.googleapis.com/v1/projects/PROJECT_ID/locations/REGISTRY_REGION/schemaRegistries/REGISTRY_NAME
    

    更改下列內容:

    • PROJECT_ID:專案 ID
    • REGISTRY_REGION:您建立結構定義登錄檔的區域
    • REGISTRY_NAME:結構定義登錄檔的名稱
  9. 點選「建立」

gcloud

如要建立 BigQuery Sink 連接器,請執行 gcloud managed-kafka connectors create 指令。


export REGISTRY_PATH=projects/PROJECT_ID/locations/REGISTRY_REGION/schemaRegistries/REGISTRY_NAME

gcloud managed-kafka connectors create CONNECTOR_NAME \
  --location=REGION \
  --connect-cluster=CONNECT_CLUSTER \
  --configs=connector.class=com.wepay.kafka.connect.bigquery.BigQuerySinkConnector,\
key.converter=org.apache.kafka.connect.storage.StringConverter,\
value.converter=io.confluent.connect.avro.AvroConverter,\
value.converter.bearer.auth.credentials.source=GCP,\
value.converter.schema.registry.url=https://managedkafka.googleapis.com/v1/$REGISTRY_PATH,\
tasks.max=3,\
project=PROJECT_ID,\
defaultDataset=DATASET_NAME,\
topics=newUsers

更改下列內容:

  • REGISTRY_REGION:您建立結構定義登錄檔的區域
  • REGISTRY_NAME:結構定義登錄檔的名稱
  • CONNECTOR_NAME:連接器名稱,例如 bigquery-connector
  • CONNECT_CLUSTER:Connect 叢集名稱
  • REGION:您建立 Connect 叢集的區域
  • PROJECT_ID:專案 ID
  • DATASET_NAME:BigQuery 資料集名稱

查看結果

如要在 BigQuery 中查看結果,請對資料表執行下列查詢。系統可能需要幾分鐘,才會將第一批訊息寫入資料表。

控制台

  1. 開啟「BigQuery」BigQuery頁面。

    前往 BigQuery 頁面

  2. 在查詢編輯器中執行下列查詢:

    SELECT * FROM `PROJECT_ID.DATASET_NAME.newUsers`
    

    請替換下列變數:

    • PROJECT_ID:您 Google Cloud專案的名稱
    • DATASET_NAME:BigQuery 資料集名稱

gcloud

使用 bq query 指令查詢資料表:

bq query --use_legacy_sql=false 'SELECT * FROM `PROJECT_ID.DATASET_NAME.newUsers`'

請替換下列變數:

  • PROJECT_ID:您 Google Cloud專案的名稱
  • DATASET_NAME:BigQuery 資料集名稱

清除所用資源

為避免因為本教學課程所用資源,導致系統向 Google Cloud 帳戶收取費用,請刪除含有相關資源的專案,或者保留專案但刪除個別資源。

控制台

  1. 刪除 Connect 叢集。

    1. 前往「Managed Service for Apache Kafka」>「Connect Clusters」(連線叢集) 頁面。

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

    2. 選取 Connect 叢集,然後按一下「刪除」

  2. 刪除 Kafka 叢集。

    1. 前往「Managed Service for Apache Kafka」>「Clusters」(叢集) 頁面。

      前往「Clusters」(叢集)

    2. 選取 Kafka 叢集,然後按一下「Delete」(刪除)

  3. 刪除 BigQuery 資料表和資料集。

    1. 前往「BigQuery」頁面

      前往 BigQuery 頁面

    2. 在「Explorer」窗格中展開專案,然後選取資料集。

    3. 展開「動作」選項,然後點按「刪除」

    4. 在「Delete dataset」(刪除資料集) 對話方塊中,在欄位輸入 delete,然後按一下「Delete」(刪除)

gcloud

  1. 如要刪除 Connect 叢集,請使用 gcloud managed-kafka connect-clusters delete 指令。

    gcloud managed-kafka connect-clusters delete CONNECT_CLUSTER \
      --location=REGION --async
    
  2. 如要刪除 Kafka 叢集,請使用 gcloud managed-kafka clusters delete 指令。

    gcloud managed-kafka clusters delete KAFKA_CLUSTER \
      --location=REGION --async
    
  3. 如要同時刪除 BigQuery 資料集和 BigQuery 資料表,請使用 bq rm 指令。

    bq rm --recursive --dataset PROJECT_ID:DATASET_NAME
    

後續步驟