更新連接器

您可以編輯連接器來更新設定,例如變更讀取或寫入的主題、修改資料轉換,或調整錯誤處理設定。

如要更新 Connect 叢集中的連結器,可以使用 Google Cloud 控制台、gcloud CLI、Managed Service for Apache Kafka 用戶端程式庫或 Managed Kafka API。您無法使用開放原始碼 Apache Kafka API 更新連接器。

事前準備

更新連接器前,請先檢查現有設定,並瞭解所做變更可能造成的影響。

更新連結器所需的角色和權限

如要取得編輯連接器所需的權限,請要求管理員在包含 Connect 叢集的專案中,授予您受管理 Kafka 連接器編輯者 (roles/managedkafka.connectorEditor) IAM 角色。如要進一步瞭解如何授予角色,請參閱「管理專案、資料夾和組織的存取權」。

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

所需權限

如要編輯連接器,必須具備下列權限:

  • 在上層 Connect 叢集上授予更新連接器權限: managedkafka.connectors.update
  • 在上層 Connect 叢集上授予清單連接器權限: 只有透過 Google Cloud 控制台更新連接器時,才需要這項權限。

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

連接器的可編輯屬性

連接器的可編輯屬性取決於其類型。以下是支援的連接器類型可編輯的屬性摘要:

MirrorMaker 2.0 來源連接器

  • 主題名稱或主題規則運算式 (以半形逗號分隔):要複製的主題。

    如要進一步瞭解這項屬性,請參閱「主題名稱」。

  • 設定:連接器的其他設定。

    如要進一步瞭解這項屬性,請參閱「 設定」。

  • 工作重新啟動政策:重新啟動失敗連接器工作的政策。

    如要進一步瞭解這項屬性,請參閱「 工作重新啟動政策」。

BigQuery Sink 連接器

  • 主題:要串流資料的 Kafka 主題。

    如要進一步瞭解這項屬性,請參閱「主題」。

  • 資料集:用於儲存資料的 BigQuery 資料集。

    如要進一步瞭解這項屬性,請參閱「資料集」。

  • 設定:連接器的其他設定。

    如要進一步瞭解屬性,請參閱「設定」。

  • 任務重新啟動政策:重新啟動失敗連接器任務的政策。

    如要進一步瞭解這項屬性,請參閱「工作重新啟動政策」。

Cloud Storage 接收器連接器

  • 主題:要串流資料的 Kafka 主題。

    如要進一步瞭解這項屬性,請參閱「主題」。

  • Cloud Storage bucket: 用於儲存資料的 Cloud Storage bucket。

    如要進一步瞭解這項屬性,請參閱「Bucket」。

  • 設定:連接器的其他設定。

    如要進一步瞭解這項屬性,請參閱「設定」。

  • 工作重新啟動政策:重新啟動失敗連接器工作的政策。

    如要進一步瞭解這項屬性,請參閱「工作重新啟動政策」。

Pub/Sub 來源連接器

  • Pub/Sub 訂閱項目:用來接收訊息的 Pub/Sub 訂閱項目。
  • Kafka 主題:要將訊息串流至的 Kafka 主題。
  • 設定:連接器的其他設定。 詳情請參閱「 設定連接器」。
  • 任務重新啟動政策:重新啟動失敗的連接器任務政策。詳情請參閱「工作重新啟動政策」。

Pub/Sub 接收器連接器

  • 主題:要從中串流訊息的 Kafka 主題。

    如要進一步瞭解該屬性,請參閱「主題」。

  • Pub/Sub 主題:要傳送訊息的 Pub/Sub 主題。

    如要進一步瞭解這項屬性,請參閱「Pub/Sub 主題」。

  • 設定:連接器的其他設定。

    如要進一步瞭解這項屬性,請參閱「設定」。

  • 任務重新啟動政策:重新啟動失敗連接器任務的政策。

    如要進一步瞭解這項屬性,請參閱「工作重新啟動政策」。

更新連接器

更新連接器時,系統會套用變更,因此資料流程可能會暫時中斷。

控制台

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

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

  2. 按一下要更新連接器的 Connect 叢集。

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

  3. 在「資源」分頁中,找出清單中的連結器,然後按一下連結器名稱。

    系統會將您重新導向至「連接器詳細資料」頁面。

  4. 按一下 [編輯]

  5. 更新連接器的必要屬性。可用屬性會因連結器類型而異。

  6. 按一下 [儲存]

gcloud

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

    啟用 Cloud Shell

    Google Cloud 主控台底部會開啟一個 Cloud Shell 工作階段,並顯示指令列提示。Cloud Shell 是已安裝 Google Cloud CLI 的殼層環境,並已針對您目前的專案設定好相關值。工作階段可能要幾秒鐘的時間才能初始化。

  2. 使用 gcloud managed-kafka connectors update 指令更新連接器:

    您可以使用 --configs 標記搭配以半形逗號分隔的鍵/值組合,或使用 --config-file 標記搭配 JSON 或 YAML 檔案的路徑,更新連線器的設定。

    以下是使用 --configs 旗標的語法,其中包含以半形逗號分隔的鍵/值組合。

    gcloud managed-kafka connectors update CONNECTOR_ID \
        --location=LOCATION \
        --connect-cluster=CONNECT_CLUSTER_ID \
        --configs=KEY1=VALUE1,KEY2=VALUE2...
    

    以下語法會搭配 --config-file 標記,以及 JSON 或 YAML 檔案的路徑。

    gcloud managed-kafka connectors update CONNECTOR_ID \
        --location=LOCATION \
        --connect-cluster=CONNECT_CLUSTER_ID \
        --config-file=PATH_TO_CONFIG_FILE
    

    更改下列內容:

    • CONNECTOR_ID:必填,要更新的連接器 ID。
    • LOCATION:必填,含有連接器的 Connect 叢集位置。
    • CONNECT_CLUSTER_ID:必填,包含連接器的 Connect 叢集 ID。
    • KEY1=VALUE1,KEY2=VALUE2...:以半形逗號分隔的設定屬性,用於更新。例如:tasks.max=2,value.converter.schemas.enable=true
    • PATH_TO_CONFIG_FILE:包含要更新設定屬性的 JSON 或 YAML 檔案路徑。例如:config.json

    使用 --configs 的範例指令:

    gcloud managed-kafka connectors update test-connector \
        --location=us-central1 \
        --connect-cluster=test-connect-cluster \
        --configs=tasks.max=2,value.converter.schemas.enable=true
    

    使用 --config-file 的範例指令。以下是名為 update_config.yaml 的範例檔案:

    tasks.max: 3
    topic: updated-test-topic
    

    以下是使用該檔案的指令範例:

    gcloud managed-kafka connectors update test-connector \
        --location=us-central1 \
        --connect-cluster=test-connect-cluster \
        --config-file=update_config.yaml
    

Go

在試用這個範例之前,請先按照「 安裝用戶端程式庫」中的 Go 設定說明操作。詳情請參閱 Managed Service for Apache Kafka Go API 參考文件

如要向 Managed Service for Apache Kafka 進行驗證,請設定應用程式預設憑證(ADC)。 詳情請參閱「為本機開發環境設定 ADC」。

import (
	"context"
	"fmt"
	"io"

	"cloud.google.com/go/managedkafka/apiv1/managedkafkapb"
	"google.golang.org/api/option"
	"google.golang.org/protobuf/types/known/fieldmaskpb"

	managedkafka "cloud.google.com/go/managedkafka/apiv1"
)

func updateConnector(w io.Writer, projectID, region, connectClusterID, connectorID string, config map[string]string, opts ...option.ClientOption) error {
	// projectID := "my-project-id"
	// region := "us-central1"
	// connectClusterID := "my-connect-cluster"
	// connectorID := "my-connector"
	// config := map[string]string{"tasks.max": "6"}
	ctx := context.Background()
	client, err := managedkafka.NewManagedKafkaConnectClient(ctx, opts...)
	if err != nil {
		return fmt.Errorf("managedkafka.NewManagedKafkaConnectClient got err: %w", err)
	}
	defer client.Close()

	connectorPath := fmt.Sprintf("projects/%s/locations/%s/connectClusters/%s/connectors/%s", projectID, region, connectClusterID, connectorID)
	connector := &managedkafkapb.Connector{
		Name:    connectorPath,
		Configs: config,
	}
	paths := []string{"configs"}
	updateMask := &fieldmaskpb.FieldMask{
		Paths: paths,
	}

	req := &managedkafkapb.UpdateConnectorRequest{
		UpdateMask: updateMask,
		Connector:  connector,
	}
	resp, err := client.UpdateConnector(ctx, req)
	if err != nil {
		return fmt.Errorf("client.UpdateConnector got err: %w", err)
	}
	fmt.Fprintf(w, "Updated connector: %#v\n", resp)
	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.Connector;
import com.google.cloud.managedkafka.v1.ConnectorName;
import com.google.cloud.managedkafka.v1.ManagedKafkaConnectClient;
import com.google.protobuf.FieldMask;
import java.io.IOException;
import java.util.HashMap;
import java.util.Map;

public class UpdateConnector {

  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 clusterId = "my-connect-cluster";
    String connectorId = "my-connector";
    // The new value for the 'tasks.max' configuration.
    String maxTasks = "5";
    updateConnector(projectId, region, clusterId, connectorId, maxTasks);
  }

  public static void updateConnector(
      String projectId, String region, String clusterId, String connectorId, String maxTasks)
      throws IOException {
    try (ManagedKafkaConnectClient managedKafkaConnectClient = ManagedKafkaConnectClient.create()) {
      Map<String, String> configMap = new HashMap<>();
      configMap.put("tasks.max", maxTasks);

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

      // The field mask specifies which fields to update. Here, we update the 'config' field.
      FieldMask updateMask = FieldMask.newBuilder().addPaths("config").build();

      // This operation is handled synchronously.
      Connector updatedConnector = managedKafkaConnectClient.updateConnector(connector, updateMask);
      System.out.printf("Updated connector: %s\n", updatedConnector.getName());
      System.out.println(updatedConnector.getAllFields());

    } catch (IOException | ApiException e) {
      System.err.printf("managedKafkaConnectClient.updateConnector 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 import managedkafka_v1
from google.cloud.managedkafka_v1.services.managed_kafka_connect import (
    ManagedKafkaConnectClient,
)
from google.cloud.managedkafka_v1.types import Connector
from google.protobuf import field_mask_pb2

# TODO(developer)
# project_id = "my-project-id"
# region = "us-central1"
# connect_cluster_id = "my-connect-cluster"
# connector_id = "my-connector"
# configs = {
#     "tasks.max": "6",
#     "value.converter.schemas.enable": "true"
# }

connect_client = ManagedKafkaConnectClient()

connector = Connector()
connector.name = connect_client.connector_path(
    project_id, region, connect_cluster_id, connector_id
)
connector.configs = configs
update_mask = field_mask_pb2.FieldMask()
update_mask.paths.append("config")

# For a list of editable fields, one can check https://cloud.google.com/managed-service-for-apache-kafka/docs/connect-cluster/update-connector#editable-properties.
request = managedkafka_v1.UpdateConnectorRequest(
    update_mask=update_mask,
    connector=connector,
)

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

You can also update the connector's task restart policy without
including the configuration, by using the `--task-restart-min-backoff`
and `--task-restart-max-backoff` flags. For example:

```sh
gcloud managed-kafka connectors update test-connector \
  --location=us-central1 \
  --connect-cluster=test-connect-cluster \
  --task-restart-min-backoff="60s" \
  --task-restart-max-backoff="90s"
Apache Kafka® 是 The Apache Software Foundation 或其關聯企業在美國與/或其他國家/地區的註冊商標。