何時應考慮導入 Apache Kafka?
當多個系統需要各自獨立消費相同事件,而且事件歷史必須保留一段時間以便日後再次讀取時,Apache Kafka 是值得考慮的分散式事件串流平台。與其只是因為需要非同步處理就選擇它,更應評估是否同時需要多重訂閱、高量處理、重新處理與容錯能力。kafka.apache.org
Kafka 常被稱為「訊息佇列」,但這個標籤本身無法完整說明它的核心價值。它擅長將系統中已發生的事實——例如訂單建立、付款完成、客戶操作或系統日誌——記錄為事件,然後讓多個應用程式與資料系統依自己的步調讀取。相較之下,若小型服務只需要一次性處理某一類背景工作,Kafka 的維運複雜度可能高於其帶來的效益。
本文先界定 Kafka 所解決的問題,再檢視會提高導入價值的訊號,以及其設計與維運上的取捨。
Kafka 是什麼類型的平台?
Kafka 圍繞用來記錄事件的主題(topic)運作。事件是代表系統中已發生事實的資料紀錄,例如「訂單已建立」、「使用者瀏覽了商品」或「感測器量測到溫度」。寫入事件的應用程式稱為生產者(producer),讀取並處理事件的應用程式則稱為消費者(consumer)。kafka.apache.org
生產者會將事件發布到主題,而消費者訂閱並讀取所需的主題。生產者不必直接知道是誰在讀取其事件。即使日後新增分析服務、通知服務或搜尋索引服務,訂單服務原則上仍可發布相同的訂單事件,而不必持續為每個系統新增獨立的整合程式碼。這就是生產者與消費者之間的鬆散耦合。kafka.apache.org
此外,Kafka 事件不會在消費者讀取後立即消失。事件會依主題層級的保留政策儲存,而消費者則管理表示其讀取進度的位置。這讓新的消費者可以從歷史紀錄開始讀取,或讓既有消費者在修正錯誤後從特定位置重新處理。kafka.apache.org
因此,將 Kafka 理解為「維護可共享的事件歷史」的平台,比將它僅理解為「傳遞訊息」的平台更有幫助。但保留不代表永久保存資料;實際保留時間仍必須依主題政策與儲存容量規劃決定。
核心元件如何協同運作?
區分 Kafka 的主要元件,有助於做出導入決策與分析事件事故。
| 元件 | 角色 | 評估是否導入時要考量的事項 |
|---|---|---|
| 主題(Topic) | 將性質相近事件分組的邏輯串流 | 需要定義事件意義、保留期限與存取權限。 |
| 分區(Partition) | 用來分割主題的有序日誌單位 | 它是吞吐量、平行度與順序保證的單位。 |
| 生產者(Producer) | 將事件寫入主題的應用程式 | 需要決定事件 key 與失敗時的重試行為。 |
| 消費者(Consumer) | 從主題讀取事件的應用程式 | 需要針對重複處理、延遲與錯誤復原進行設計。 |
| 消費者群組(Consumer group) | 分攤工作的消費者集合 | 同一群組內的消費者會分配各個分區。 |
| Broker | 儲存與提供事件的 Kafka 伺服器 | 它是複寫、故障網域與儲存容量的維運單位。 |
一個主題會分成一個或多個分區。分區是有序的事件日誌,Kafka 透過多個分區來平行化讀寫。因此,分區數量不只是設定值,而是同時反映目標吞吐量、消費者平行度與順序需求的設計決策。kafka.apache.org
消費者群組是執行相同工作的一組消費者執行個體。例如,若多個消費者執行個體都將訂單事件載入資料倉儲,它們可以組成一個群組。在群組內,每個分區可指派給一個消費者,以分攤處理負載。相對地,通知服務與分析服務屬於不同群組,因此兩者都能獨立讀取相同的訂單事件。kafka.apache.org
此結構有利於擴展,但群組中活躍消費者的數量多於分區,並不表示所有消費者都能同時處理更多分區。不能只靠增加執行個體數量,就期待平行度無限制成長。從一開始,分區規劃就必須同時考量實際 key 分布與未來擴展需求。
哪些問題會讓 Kafka 更適合?
最強的導入訊號,是多個系統需要以不同用途、不同速度消費同一事件的情況。重點不在消費者數量本身,而在於消費者是否需要能獨立於生產者演進。
以電子商務系統建立訂單為例。一開始可能只需更新訂單資料庫即可。之後可能加入庫存保留、付款流程、客戶通知、詐欺偵測、搜尋與推薦資料更新,以及分析資料載入。若每項功能都持續透過同步呼叫連接至訂單服務,其中任一功能的延遲或失敗都可能影響訂單處理路徑,整合關係也會變得複雜。
在這種情況下,訂單服務可發布 order created 事件,而各下游系統透過各自獨立的消費者群組讀取所需事件。Kafka 的一個關鍵使用案例,是能在不直接修改既有生產者的情況下新增消費者。kafka.apache.org
以下情況尤其值得評估:
- 有訂單、付款或會員狀態變更等多個業務系統都會參照的核心事件。
- 使用者點擊、頁面瀏覽、營運日誌或量測值等資料持續累積。
- 分析、通知、索引與資料倉儲載入都需要同一份來源事件。
- 即使消費者的處理速度不同,或故障後的復原點不同,生產流程仍必須持續進行。
- 出現新的使用案例時,來源服務若要直接連接每一個下游系統會很繁瑣。
其中任一條件不一定代表必須使用 Kafka。但若同時符合多項條件,且各資料流很可能持續成長,事件串流架構可能比簡單的點對點整合帶來更大效益。
Kafka 如何吸收高流量與劇烈波動?
Kafka 的設計會透過分區分散事件讀寫,因此可用於持續產生大量事件的資料流。常見例子包括日誌彙整、使用者活動追蹤、監控指標、IoT 量測資料與交易事件。kafka.apache.org
Kafka 在此的作用,是放寬生產與消費速度必須始終一致的耦合。例如,若特定時段事件量激增,消費者可能無法立即處理全部事件。只要事件被保留,消費者就能稍後追上累積的待處理資料。這讓消費者可以獨立調整處理速率,而不會阻塞生產者。
這並不代表延遲會消失;而是代表延遲能以已記錄事件的累積待處理量來管理。消費者延遲(consumer lag)是一項維運指標,用來顯示消費者落後最新事件的程度。若延遲持續增加,應調查消費者效能、外部相依性、分區分布與錯誤重試。Kafka 提供以 JMX 為基礎的監控指標,而正式環境也必須考量監控存取路徑的安全性。kafka.apache.org
評估吞吐量需求時,與其籠統地說「流量很高」,不如分別回答下列問題:
- 每秒或每個時段會產生多少事件?
- 單一事件的平均與最大大小是多少?
- 尖峰會持續多久?
- 可接受多少消費者延遲?
- 故障後必須在多快時間內處理完待處理資料?
- 事件必須保留多久?
回答這些問題後,可以看出分區數量、儲存容量、複寫、消費者擴展與重新處理時間是彼此相連的議題。Kafka 提供高吞吐量的基礎,但實際效能與成本會因事件大小、key 偏斜、保留政策與消費者邏輯瓶頸而不同。
為何重新處理是導入 Kafka 的重要理由?
即時處理是事件到達後立即產生結果的工作,例如訂單建立後更新庫存、偵測符合特定條件的交易,或彙總分鐘等級的指標。但重新處理歷史資料的需求,可能和即時處理一樣重要。
重新處理有許多原因。修正消費者程式碼中的錯誤後,可以重建遺漏或計算錯誤的結果。導入新的分析規則時,可從既有事件歷史建立衍生資料。若因故障導致消費中斷,也可從最後已處理的位置再次讀取來復原。在 Kafka 中,事件不會在消費後立即移除,並可在保留政策允許的範圍內重新讀取。kafka.apache.org
例如,假設客戶行為事件最初僅用於彙總每日訪客數。若之後需要依獲客管道進行轉換分析,只要事件包含必要欄位且保留期限仍有效,另一個消費者群組就能讀取歷史事件並產生新的分析結果。這項工作可以獨立進行,不必停止既有的彙總消費者,也不必對來源服務的資料庫執行大規模查詢。
然而,能夠重新處理並不會自行解決資料品質問題。若事件缺少必要識別碼、發生時間戳記或版本資訊,或結構描述的意義在未做相容性管理下改變,即使能讀取歷史資料,也很難產生可靠結果。此外,若需求是重新處理早於保留期限的資料,僅靠 Kafka 主題可能無法滿足。因此,若重新處理是導入理由,應先確定「要重播什麼、重播多久,以及其意義是什麼」。
Kafka 也能連接資料庫與外部系統嗎?
Kafka 不僅可用於在服務之間傳遞事件,也可作為資料管線的中心流。變更資料擷取(CDC)是一種將資料庫中發生的變更送入資料流的方法,當營運資料的變更必須反映在分析、搜尋或其他服務時,可以考慮使用。Kafka Connect 為與外部系統的定期資料輸入與輸出整合,提供 API 與 connector 模型。kafka.apache.org
此組態可能適用的例子包括:
- 持續將營運資料庫的變更傳送至分析儲存庫。
- 將多個應用程式的日誌與指標收集到共享資料流中。
- 將一個系統產生的資料反映到另一個儲存庫中的索引或衍生資料表。
- 在地端與雲端環境之間建立持續性的資料流。
使用 connector 並不會消除資料模型差異、刪除語意、順序問題、存取管理或目標系統寫入限制。特別是將資料庫變更當作事件時,必須區分「某一列已變更」這項事實與「訂單已確認」這項業務事件。前者較接近儲存層變更,後者則是具有領域意義的業務事件。將兩者視為相同,可能使消費者過度依賴儲存結構。
因此,將 Kafka 導入資料管線時,若不僅能減少連線數量,還能釐清資料擁有者、結構描述與變更責任,可靠性會更高。
順序可保證到什麼程度?為何 key 設計重要?
在 Kafka 中,事件順序只在分區內保證,而非在整個主題中保證。多個分區可以實現平行處理,但它們之間不存在單一的全域順序。kafka.apache.org
例如,若訂單狀態必須依 created、payment completed 與 shipping started 的順序處理,可使用訂單 ID 作為 key,讓同一訂單的事件記錄到相同分區。如此便能以該訂單為單位使用紀錄順序。客戶狀態變更也可採用類似設計,將客戶 ID 作為 key。
反過來說,若所有訂單事件都必須依整體時間順序逐一處理,實際上可能需要接近單一分區的選擇。在這種情況下,順序可能較單純,但平行處理能力會受到限制。全域排序與高平行度並不是可以無限制同時取得的特性。
key 的選擇還有另一個問題。若特定客戶或裝置產生異常大量事件,該 key 可能集中於一個分區。這可視為 key 偏斜,而只有部分消費者可能變得過度忙碌。因此,key 應代表需要順序的業務單位,同時也應依預期資料分布評估,確保不會產生過度偏斜。
記錄順序需求時,不應只停留在「順序很重要」。最好具體化如下:
- 在什麼識別碼範圍內需要保證順序?
- 需要的是事件發生時間順序,還是紀錄寫入順序?
- 要如何處理延遲到達的事件?
- 若事件順序錯亂,會造成什麼業務錯誤?
- 即使犧牲平行度,是否仍需要全域排序?
這些答案決定主題分離方式、key、分區數量與消費者邏輯。
應如何理解重複處理與恰好一次處理?
考量故障與重試時,Kafka 消費者需要以預設至少一次處理(at-least-once processing)為前提進行設計。例如,若消費者完成事件處理後,在記錄處理位置之前停止,復原後可能再次讀取同一事件。因此,同一事件有可能被重複處理。kafka.apache.org
實務上的解法是使消費者邏輯具備冪等性(idempotency)。冪等性是指即使多次執行相同操作,最終結果仍相同的特性。例如,像 set the status of order 123 to delivered 這樣的操作,可以設計為重複執行相同狀態更新也不會實質改變結果。相對地,unconditionally adds 1,000 points 這類操作若收到同一事件兩次,可能產生不同結果,因此需要去重策略,例如記錄事件 ID,或在目標儲存庫使用唯一性限制。
當讀取、處理與寫入皆在 Kafka 主題內串接時,Kafka 可透過交易與 read_committed 隔離等級支援恰好一次處理組態。不過,這不應被理解為每種外部影響都會自動只發生一次。Kafka 以外的副作用,例如外部資料庫更新、電子郵件寄送或付款 API 呼叫,仍需要與目標系統協調並另行設計。kafka.apache.org
因此,導入前應針對每個消費者詢問以下問題:
- 同一事件被處理兩次時會發生什麼事?
- 每個事件是否都有可用於偵測重複的 ID?
- 結果儲存庫是否能防止重複,或支援安全更新?
- 外部呼叫失敗或回應不明確時,重試標準為何?
- 重新處理時,要如何處理已執行的副作用?
若未回答這些問題就導入 Kafka,傳輸本身可能很可靠,但重複或不一致的業務結果仍會難以察覺。
容錯與耐久性會自動得到保證嗎?
Kafka 可透過複寫主題分區來因應 broker 故障。對於重視分區複本、broker 故障時持續運作,以及讓眾多消費者分散負載的資料流而言,這可能是重要優勢。kafka.apache.org
不過,「因為使用 Kafka,資料永遠不會遺失」這個結論並不正確。實際的耐久性與可用性取決於複寫因子、生產者確認設定、可能同時發生的故障範圍、保留政策及維運程序。即使有複本,若複本位於相同故障網域、重要設定未達需求等級,或維運人員未驗證復原程序,結果仍可能不如預期。kafka.apache.org
將容錯需求明確寫下來會有所幫助。例如:「若一個 broker 停機,訂單事件的生產與消費必須持續」、「消費者故障後可接受重複,但不可遺漏」或「指定期間內的事件必須能重新處理」。這些需求不僅決定複寫與確認設定,也決定消費者冪等性、監控、儲存容量與復原演練。
由於可重新處理的歷史資料可能成為重要資料資產,應另外評估主題是否包含個人資訊或敏感業務資料。存取控制與維運介面安全性並非可事後處理、且與資料流設計分離的工作。Kafka 維運也需要針對管理存取進行安全設定,包括監控。kafka.apache.org
Kafka 是否總是比簡易工作佇列或同步 API 更好?
不是。Kafka 並非每一種非同步需求的自動替代方案。若需求更接近「將一張圖片轉換一次」、「產生報表後只回傳結果」,或「讓單一消費者取得並處理工作」,而長期保留、多重訂閱與重新處理並非核心需求,則較簡單的工作佇列或受管理服務可能在成本與維運負擔上更適合。Kafka 的主要優勢會在大型事件流、多個獨立消費者與重複利用保留歷史三者結合時顯現。kafka.apache.org
同步 API 也有不同的角色。使用者點擊按鈕後需要立即得知成功或失敗結果的請求,自然適合 request-response API。該請求完成後,將這項事實通知下游系統的流程,可以拆分為事件。換言之,與其只選擇同步呼叫或 Kafka 其中之一,通常更適合以 API 處理使用者互動,並以事件處理下游的非同步扇出。
以下比較可簡化決策:
| 主要需求 | 優先評估的方法 | Kafka 特別具優勢的條件 |
|---|---|---|
| 一次處理一項工作 | 簡易工作佇列或受管理的非同步服務 | 多個獨立系統必須讀取相同工作結果或事件時 |
| 需要立即結果的請求 | 同步 API | 請求完成後,許多不同的下游工作必須以非同步方式扇出時 |
| 在系統間傳輸資料 | 直接整合或檔案/批次方式 | 持續資料流、多個目的地與重新處理需求同時存在時 |
| 收集日誌、行為或量測資料 | 收集工具與儲存空間 | 多個消費者必須獨立處理高量串流時 |
| 管理狀態變更歷史 | 業務資料庫 | 必須重播事件以重建狀態或衍生資料時 |
此表並非絕對的產品選擇規則。團隊既有的平台、受管理服務的可用性、安全政策及維運人力也會影響決策。關鍵在於要解決的資料流性質,而不是功能清單。
維運與治理必須準備什麼?
導入 Kafka 不只是新增一個應用程式函式庫,還需要一套持續管理主題、分區、複寫、保留、存取權限、監控與容量的維運模式。Kafka 提供 JMX 指標,但只有在確定哪些指標會觸發警示、由誰回應,以及如何復原後,維運資訊才能產生真正價值。kafka.apache.org
首先,必須管理事件契約。事件契約不僅包含欄位名稱與資料型別,還包括每個欄位的業務意義、是否為選填、版本變更的處理方式,以及生產時間與發生時間的區別。需要有相容性標準,避免生產者刪除欄位或變更欄位意義時,消費者在未察覺下錯誤運作。
接著,主題政策必須明確。對每個主題,都需要決定以下事項:
- 包含哪些事件,以及誰擁有這些事件。
- 保留期限與儲存容量的判定標準。
- 哪些順序與吞吐量需求決定分區數量與 key。
- 複寫與生產者確認設定預期達成何種故障承受等級。
- 誰可以生產與消費,以及如何保護敏感資料。
- 消費者延遲到達何種程度時開始調查與回應。
容量規劃同樣重要。較長的保留期限或更多複本會提高儲存需求。若消費者停止很長一段時間後仍必須能重新處理,便可能需要相應保留歷史資料。相反地,較短保留期限可降低成本,但會限制可用於故障復原或新增消費者的歷史資料範圍。這項選擇定義的不只是成本,也包括產品能力與可復原性的範圍。
在維運責任不明確的組織中,共享 Kafka 平台反而可能增加相依性問題。針對哪些變更與事故分別屬於主題擁有者、平台維運人員、安全負責人與消費者開發團隊的責任達成共識,與技術組態同樣重要。
導入前應用哪些問題來做決定?
最能區分是否該導入 Kafka 的問題,不是「我們需要非同步訊息嗎?」更準確的問題是:**多個獨立消費者是否需要持續讀取大規模事件歷史,並在發生延遲或故障後重新處理?**若答案明確為是,代表需求很可能符合 Kafka 的核心特性。kafka.apache.orgkafka.apache.org
開始討論導入時,可使用以下檢查清單:
- 多個消費者:目前或近期內,是否有多個系統需要獨立使用同一事件?
- 歷史價值:事件是否必須在消費後保留,並可為了修正錯誤、稽核或新分析而再次讀取?
- 處理規模:持續高量的寫入或尖峰流量,是否需要將生產與消費解耦?
- 順序範圍:問題是否可透過客戶或訂單等 key 的順序解決,而非要求全域排序?
- 重複處理:每個消費者能否安全處理或識別重複事件?
- 契約管理:是否有負責人與流程可管理事件結構描述與意義的變更?
- 維運準備度:是否有負責人能觀察並回應延遲、儲存容量、broker 故障、權限與重新處理?
- 替代方案比較:需求是否能僅透過單一消費者的工作分派或 request-response,以更簡單的方式滿足?
導入 Kafka 時,不必從一開始就讓每一項都完美。但若第 1 至第 5 項的需求很強,卻缺乏第 6 與第 7 項的準備,技術可行性與可維運系統之間可能存在很大落差。可先在一條小範圍資料流中驗證事件契約、重複處理、延遲觀察與重新處理。
結論:必須共享事件歷史時,Kafka 的威力最強
Apache Kafka 不只是非同步移動訊息的工具;它是一個保留多個系統共享之事件流,並讓各系統獨立消費這些資料流的平台。在需要對相同事件進行多重訂閱、平行處理高量資料流、在延遲後追上進度,以及重新處理歷史紀錄的環境中,它的導入價值會提高。kafka.apache.orgkafka.apache.org
相反地,對於將一次性工作交給單一消費者、以立即回應為核心的請求,或必須盡量降低維運負擔的小型流程,較簡單的替代方案可能更合適。選擇 Kafka 時,除了吞吐量外,也應評估是否已準備好管理分區層級的順序、重複處理、保留政策、事件契約、安全性與可觀測性。這些條件越完備,Kafka 就越能成為降低服務間耦合、擴展資料使用方式的基礎。