連接器總覽

下表列出 Managed Service for Apache Kafka 支援的 Kafka Connect 連接器類型。您可以使用這些連接器,將 Apache Kafka 與應用程式和其他 Google Cloud 服務整合。

連接器 說明 用途
MirrorMaker 2.0 將某個 Kafka 叢集的主題和資料,複製到另一個 Kafka 叢集。 資料複製、災難復原、資料遷移
BigQuery Sink 將 Kafka 主題的資料串流至 BigQuery 資料表。 資料倉儲、分析
Cloud Storage Sink 將 Kafka 主題的資料串流至 Cloud Storage bucket。 資料湖泊擷取、資料封存
Pub/Sub Sink 將 Kafka 主題的資料串流至 Pub/Sub 主題。 服務整合、即時通知
Pub/Sub 來源 將 Pub/Sub 訂閱項目的訊息串流至 Kafka 主題。 即時資料擷取、事件導向架構

轉換者

轉換器負責 Kafka 記錄資料的序列化和還原序列化作業。這些轉換器會在 Kafka 主題的原始位元組格式,以及 Kafka Connect 使用的內部結構化資料表示法之間進行轉換。

  • 對於接收器連接器,轉換器會將主題的線路格式資料還原序列化為 Kafka Connect 內部資料格式,連接器會使用該格式寫入目標系統。

  • 對於來源連接器,轉換器會將 Kafka Connect 內部資料格式的資料序列化為 Kafka 主題的指定線路格式。

轉換器可確保連接器讀取或寫入 Kafka 記錄時,採用與外部系統相容的格式。

設定連接器時,請設定下列屬性:

  • 鍵轉換器 (key.converter):用於序列化和還原序列化 Kafka 記錄鍵的轉換器。

  • 值轉換器 (value.converter):用於序列化和還原序列化 Kafka 記錄值的轉換器。

如未指定轉換器,預設轉換器類型為 org.apache.kafka.connect.converters.ByteArrayConverter,會以原始位元組格式傳遞資料。

支援的轉換器

Managed Service for Apache Kafka 支援下列內建轉換器:

轉換者格式
io.confluent.connect.avro.AvroConverter Apache Avro
org.apache.kafka.connect.converters.BooleanConverter 布林值
org.apache.kafka.connect.converters.ByteArrayConverter

位元組陣列

預設轉換器類型。保留兩個系統中訊息的確切內容。

org.apache.kafka.connect.converters.DoubleConverter 雙精度值
org.apache.kafka.connect.converters.FloatConverter 浮點值
org.apache.kafka.connect.converters.IntegerConverter 整數
org.apache.kafka.connect.json.JsonConverter

JSON

如果 JSON 資料沒有結構定義,也請設定 value.converter.schemas.enable=false

org.apache.kafka.connect.converters.LongConverter Long
org.apache.kafka.connect.converters.ShortConverter 短文案
org.apache.kafka.connect.storage.StringConverter 字串

轉換器的選擇取決於連接器類型,以及您在 Kafka 中儲存的資料。詳情請參閱特定連結器的說明文件。

Tasks

連接器會建立一或多個工作,並以平行方式運作,藉此轉移資料。如要設定連接器建立的工作數量上限,請設定連接器的 tasks.max 設定屬性。連接器建立的任務數量可能少於這個值。

增加 tasks.max 的值可提高總處理量,但也會增加資源消耗 (CPU 和記憶體)。最佳值取決於工作負載,以及分配給 Connect 叢集工作站的資源。對於接收器連接器,Kafka 主題分區數量也會影響平行處理。

任務重新啟動政策

您可以設定連接器的工作重新啟動政策,決定發生失敗時的行為。連結器支援下列政策:

  • 請勿重新啟動。連接器不會重新啟動失敗的工作。這項政策是預設行為。這項功能有助於偵錯,或在發生錯誤後需要手動介入的情況下使用。

  • 以指數輪詢方式重新啟動。連接器會在延遲一段時間 (稱為「退避」期間) 後,重新啟動失敗的任務。每次後續失敗時,延遲時間會呈指數增加。建議大多數正式環境工作負載採用這項政策。

    如果使用指數輪詢政策,請一併設定最短和最長輪詢時間。輪詢持續時間下限應大於 60 秒,輪詢持續時間上限應小於 7200 秒。

轉換和述詞

Managed Service for Apache Kafka 支援預設的 Kafka Connect 轉換述詞

轉換功能可讓您在個別訊息傳送至 Managed Service for Apache Kafka (適用於來源連接器) 或外部系統 (適用於接收器連接器) 前修改訊息。您可能會使用轉換來遮蓋私密資料、新增時間戳記或重新命名欄位。

述詞可讓您根據特定條件篩選資料,並根據訊息屬性判斷轉換指令適用的訊息。

舉例來說,如要設定接收器連接器,使其忽略含有 DoNotProcess 標頭鍵的訊息,請新增下列設定:

transforms=dropMessage
transforms.dropMessage.type=org.apache.kafka.connect.transforms.Filter
transforms.dropMessage.predicate=hasKey
predicates=hasKey
predicates.hasKey.type=org.apache.kafka.connect.transforms.predicates.HasHeaderKey
predicates.hasKey.name=DoNotProcess

這項設定會執行下列操作:

  1. 設定名為 hasKey 的述詞,類型為 org.apache.kafka.connect.transforms.predicates.HasHeaderKey。這項述詞會比對含有以 DoNotProcess 為鍵的標頭的所有郵件。

  2. 設定名為 dropMessage 的轉換,類型為 org.apache.kafka.connect.transforms.Filter。這項轉換會捨棄符合所設定述詞的所有訊息。

  3. 將轉換連結至述詞 hasKey。這可確保只有含有 DoNotProcess 標頭鍵的訊息會遭到轉換作業捨棄。

後續步驟

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