Kafka Connect:連接器的執行時
· tech
📑 目錄
收尾 stateless 這一批的最後一個:Kafka Connect。它是一套專門在 Kafka 與外部系統(資料庫、S3、Elasticsearch)之間搬資料的框架——你不用每次都手寫消費者/生產者,只要配一份 connector。從 infra 角度,它是這個系列「無狀態運算 + 借外部狀態」模式最漂亮的一個示範,因為它把自己的狀態,存回了它所依賴的 Kafka。
拓撲與狀態:worker 無狀態,狀態全存回 Kafka
這就是 Connect 最聰明的 infra 設計:它不自己維護狀態,而是把 Kafka —— 它本來就依賴的東西 —— 當成自己的 state store。 worker 因此變成純粹的無狀態運算層,可拋、可換、可隨便加。這跟 Spark 借 S3、Airflow 借 metadata DB 是同一個配方,只是 Connect 借的,是它自己腳下的 Kafka。
source 與 sink:兩個方向,平行度各有天花板
Connect 做的事對稱得很乾淨:source 把外部資料灌進 Kafka、sink 把 Kafka 的資料倒去外部。而它能開多平行,由 tasks.max 設上限,但真正的天花板在兩端:
tasks.max 設平行度上限,但真正搬得多快,卡在兩端的分片:source 受來源的分片數(例如幾張表、DB 的幾個 partition)限;而 sink task 本質就是一個 consumer group 的成員,所以它的平行度受 topic partition 數限——partition 有幾個,sink 最多就幾個 task 有效。加 worker 加不過這道牆HA、擴展、監控、在 k8s 上
- 擴展與 rebalance:加 worker,task 就會在 worker 之間重新分配(rebalance),行為跟 consumer group 一模一樣。舊版 Connect 的 rebalance 是「stop-the-world」——一動全停;新版的 incremental cooperative rebalancing 只搬需要動的 task,大幅減少抖動。
- HA:distributed mode 下,workers 組成一個群組、透過 Kafka 協調,任何 worker 掛掉,它的 task 自動 rebalance 到存活的 worker——因為狀態在 Kafka,這是無痛的。唯一的單點是 Kafka 叢集本身:它不健康,Connect 就動不了。
- 監控與故障:盯 connector / task 的 status(
RUNNING/FAILED/PAUSED)、source/sink 的 lag、throughput。有個大坑:taskFAILED後預設不會自動重啟,得靠 REST APIrestart或外部監控補上;遇到解不了的「毒訊息」(schema 不符),設 dead letter queue 把它丟到一邊、別卡住整條。 - 在 k8s 上:worker 無狀態,所以它是這批工具裡最好上 k8s 的——就是一個 stateless 的 Deployment,甚至能直接 autoscale;connector 全靠 REST API(
POST /connectors)管理,社群常用 Strimzi operator 把這些兜起來。
反思
Connect 把「借外部狀態」玩到了極致
寫完這批 stateless 工具,Kafka Connect 給我的收尾特別漂亮:它不只是「不自己存狀態」,而是把自己依賴的 Kafka,直接拿來當狀態儲存。Spark 借 S3、Airflow 借 metadata DB,而 Connect 借的是它腳下那個 Kafka——省得再多養一個狀態儲存。這讓我把這批工具的共通配方看得更清楚了:無狀態運算層 = 純運算的 worker + 一個被指定來扛狀態的外部儲存。 認出這個配方,你看任何號稱「無狀態、好水平擴」的服務,都會立刻去問同一句:它的狀態,寄放在誰那裡? 找到那個「誰」,你就找到了它真正的命門——對 Connect 來說,就是 Kafka。
rebalance 是「工作單位可自由搬動」的帳單
Connect 好擴、掛了好修,根源都是同一件事:task 是可以在 worker 之間自由重新分配的。但這份自由不是免費的——每次 worker 進出,都要 rebalance,而 rebalance 本身會讓正在搬的資料短暫停頓。這跟 consumer group 的 rebalance、跟 K8s 把 Pod 重排,是同一種取捨:你讓工作單位變得可以隨處搬動,換來了彈性與韌性,但搬動的那一刻要付停頓的帳。 成熟的系統不是消滅這個帳單(消不掉),而是把它壓到最小——Connect 的 incremental cooperative rebalance、K8s 的 PodDisruptionBudget,都是在做同一件事:讓「搬動」盡量不驚動還在好好幹活的部分。
好基礎設施的樣子:把重複的事變成一份設定
Connect 最打動我的,是它把「在 Kafka 和外部系統之間搬資料」這件每個團隊都要做一百遍的事,變成了填一份 JSON、POST 給一個 API,而不是每次都手寫一個消費者、一個生產者、自己管 offset。這種「把 80% 的常見需求收斂成宣告式設定,只有真的特殊才寫程式」的設計,是我心中成熟基礎設施的共通長相——K8s 的宣告式、Airflow 的 operator,都是這個精神。它把工程師從重複勞動裡解放出來,去做真正需要判斷的事。這也正好收束了 stateless 這一批:Spark、Airflow、Connect 三個無狀態運算層,骨架其實一致——無狀態 worker、借外部狀態、可彈性擴。下一篇,就把這些 stateless 的、跟前面那些 stateful 的,兜成一個完整的資料平台。