下表列出 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 資料沒有結構定義,也請設定 |
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
這項設定會執行下列操作:
設定名為
hasKey的述詞,類型為org.apache.kafka.connect.transforms.predicates.HasHeaderKey。這項述詞會比對含有以DoNotProcess為鍵的標頭的所有郵件。設定名為
dropMessage的轉換,類型為org.apache.kafka.connect.transforms.Filter。這項轉換會捨棄符合所設定述詞的所有訊息。將轉換連結至述詞
hasKey。這可確保只有含有DoNotProcess標頭鍵的訊息會遭到轉換作業捨棄。