iT邦幫忙

2026 iThome 鐵人賽

DAY 24
0
Software Development

遠古聖遺物改造工程:遺留系統全面重構實務指南系列 第 24

[Day 24] 分散式服務之間的資料要如何同步?

  • 分享至 

  • xImage
  •  

分散式服務之間的資料要如何同步?

前一章確認分散式架構會增加跨服務通訊、部分失敗與資料一致性的處理成本。當一項功能需要多個服務共同完成時,這些服務不能直接修改彼此擁有的保存內容,只能透過已定義的互動契約交換要求或已發生的結果。

「同步」在這個情境中可能代表兩件不同的事。同步呼叫表示呼叫端會等待回應,資料同步則表示不同服務保存的狀態要在允許時間內達到已定義的一致結果。採用同步呼叫不代表多個服務會自動形成同一筆交易,採用非同步訊息也不代表資料可以任意延遲。設計前要先確認資料擁有者、允許延遲與失敗後狀態,再選擇通訊方式及工具。

先定義需要同步的結果

直接要求「服務之間的資料要一致」仍不足以形成設計。每項跨服務資料至少要回答下列問題:

  • 指定哪一個服務負責接受變更並判斷功能規則,其他服務不得同時成為相同資料的寫入來源。
  • 說明其他服務需要完整內容、部分欄位、彙總結果,或只需要知道某件事已經發生。
  • 定義變更後多久內必須可見,以及延遲期間可以顯示既有結果、待處理狀態或停止相關功能。
  • 定義訊息遺失、重複、延遲及順序錯置時,對外結果與修正方式。
  • 指定哪些操作需要立即確認成功,哪些工作可以稍後完成。
  • 記錄資料保存期限、重播範圍、追蹤需求與可以接受的人工處理時間。

如果兩個服務都能直接修改同一項狀態,就要另外處理衝突、覆寫與循環同步。較清楚的設計是保留單一寫入擁有者,其他服務按照需要保存唯讀副本或衍生結果。這些副本可以暫時落後,但延遲上限與修正流程必須明確。

同步呼叫與非同步訊息傳遞

同步應用程式介面(Application Programming Interface, API)呼叫適合呼叫端必須立即取得結果,且後續步驟無法在結果未知時繼續的情境。非同步訊息傳遞則先接受工作或記錄事件,再由其他服務獨立處理,適合允許延遲、需要分開處理時間或有多個獨立接收者的情境。

比較面向 同步 API 呼叫 非同步訊息傳遞
回應方式 呼叫端在期限內等待成功、失敗或逾時結果 發送端確認訊息已被接受後即可繼續,處理結果由後續狀態或另一則訊息表示
時間相依 呼叫端與接收端必須在同一段操作期間可用 接收端可以稍後處理,但訊息代理與積壓容量要能支撐延遲
失敗狀態 逾時時可能無法判斷接收端是否已完成,重試前要確認冪等規則 訊息可能重複、延遲或無法處理,要設計確認、重試與死信流程
資料可見性 成功回應可以立即帶回結果,但不保證其他服務的副本已經更新 各接收端按照自己的進度更新,通常形成最終一致性(Eventual Consistency)
接收者數量 發送端通常明確知道要呼叫的介面 可以讓多個獨立接收者訂閱相同事件
處理量調整 呼叫端與接收端的容量會直接互相影響 消費者可以按照積壓量調整處理速度,但必須限制佇列成長
問題追查 要串連請求、逾時、重試與接收端結果 要串連發佈、路由、消費、重試、死信與最終套用結果

同一個跨服務流程可以同時使用兩種方式。例如,建立操作需要立即確認擁有者是否接受,可以使用同步呼叫。接受後產生的狀態變更通知可以透過事件傳給其他服務。選擇基準是每一步的回應與一致性需求,不需要為整個系統固定一種方式。

理解訊息傳遞的基本角色

使用非同步方式時,訊息代理(Message Broker)位於生產者與消費者之間,負責接收、保存或路由訊息。它降低雙方必須同時可用的時間相依,但不會替程式決定訊息內容、處理順序或資料正確性。

角色或構成 責任
生產者(Producer) 建立訊息、指定目的地或路由資訊,並確認訊息是否已由代理接受
訊息代理 按照佇列、主題、交換器或其他規則保存與傳遞訊息
消費者(Consumer) 讀取訊息、驗證結構、執行本身擁有的狀態變更,並回報處理進度
佇列(Queue) 保存待處理訊息,常用於讓多個消費者共同分派工作
主題(Topic) 依名稱彙整一類事件,讓不同訂閱者獨立讀取
消費者群組(Consumer Group) 讓同一類消費者共同分擔主題的分割區,同一事件仍可由其他群組獨立處理

命令與事件也要分開命名。命令要求特定責任執行一項動作,可能被拒絕,例如「要求重新計算」。事件記錄已經發生且不能由消費者撤銷的事實,例如「計算結果已更新」。如果把事件寫成命令,接收者會誤以為可以改變發送端狀態。如果把命令寫成事件,則可能失去拒絕、逾時與結果回覆的語意。

RabbitMQ 如何分派與路由訊息

在 RabbitMQ 常見的 AMQP 0-9-1 模型中,生產者將訊息發佈到交換器(Exchange),交換器再按照類型、繫結(Binding)與路由鍵將訊息送往一個或多個佇列。直接交換器可以按照完整路由鍵分派,主題交換器可以按照路由鍵模式分派,扇出交換器則會將訊息送往所有已繫結的目的地。RabbitMQ 的交換器說明將交換器定位為代理內的路由機制,佇列則保存等待消費的訊息。

這種模型適合下列情境:

  • 多個工作執行者要從同一個佇列分擔命令或背景工作。
  • 訊息要依種類、來源或其他路由條件送往不同佇列。
  • 每個接收責任需要自己的佇列、重試規則與失敗處理方式。
  • 訊息完成處理後即可移除,不需要長期保留完整事件歷程。

可靠傳遞還要同時處理發佈端確認與消費端確認。發佈端確認只表示 RabbitMQ 已按照設定接手訊息,不代表消費者已完成工作。消費端應該在本身的狀態變更成功後才確認訊息,否則處理失敗時可能無法再次取得。RabbitMQ 的可靠性說明也指出,未收到發佈確認而重新傳送時可能產生重複訊息,因此消費者仍要具備冪等性。

Apache Kafka 如何保留與重播事件

Apache Kafka 將事件追加到主題的分割區(Partition)。消費者依照位移量(Offset)記錄每個分割區已讀取的位置,主題則按照保留時間或容量規則保存事件。新的消費者群組可以從指定位置讀取,既有群組也可以調整位移量後重新處理仍在保留範圍內的事件。

同一消費者群組內,一個分割區在同一時間會分配給一個消費者。不同群組則可以各自讀取相同主題。Apache Kafka 的核心介紹說明事件會依鍵值進入分割區,順序保證只存在於同一個主題分割區內,主題保留的事件可以按照需要再次讀取。

這種模型適合下列情境:

  • 多個獨立接收責任需要各自讀取同一份事件歷程。
  • 需要在保留期限內重播事件,以重建衍生狀態或修正消費程式。
  • 事件量與持續寫入量需要透過多個分割區分散處理。
  • 資料串流處理需要按照同一鍵值維持局部順序。

增加分割區可以提高平行處理能力,也會限制整體順序。需要依序處理的事件必須使用穩定鍵值進入同一分割區。分割區數量、鍵值分布與消費者數量都要以實際事件分布驗證,不能只增加數量就假設處理能力會等比例提高。

按照需求選擇 RabbitMQ 或 Apache Kafka

RabbitMQ 與 Apache Kafka 都能傳遞訊息,但常用模型與主要取捨不同。比較時應以準備採用的佇列類型、保留設定、確認方式與用戶端行為為準,不能只比較產品名稱。

比較面向 RabbitMQ 的交換器與佇列模型 Apache Kafka 的主題與分割區模型
主要處理方式 將訊息路由到佇列,由消費者取得並確認處理 將事件追加到分割區,由消費者按照位移量讀取
路由 交換器與繫結可以在代理端進行直接、模式或扇出路由 生產者以主題與鍵值決定目的地,消費者通常訂閱完整主題
工作分派 同一佇列的多個消費者可以共同分擔訊息 同一群組的消費者共同分擔主題分割區
多個接收責任 每個責任通常建立獨立佇列 每個責任通常建立獨立消費者群組
保留與重播 傳統佇列通常在確認後移除訊息,是否需要串流或額外保存要另行設計 事件在保留範圍內不因讀取而刪除,可以調整位移量重播
順序 要依佇列類型、路由方式、消費者數量與重新排隊行為驗證 保證同一分割區內的順序,不保證不同分割區的整體順序
常見用途 工作佇列、命令分派、需要彈性代理端路由的訊息 事件歷程、資料串流、需要多個群組獨立處理或重播的事件
主要維運項目 交換器與佇列拓撲、無法路由訊息、確認、積壓與死信 主題與分割區、保留容量、鍵值分布、消費者延遲與群組重新分配

選擇工具前,可以用相同事件資料與故障條件完成一個概念驗證(Proof of Concept, POC):

  1. 驗證正常發佈、路由、消費與多個接收責任的結果。
  2. 中斷生產者、訊息代理與消費者,確認恢復後是否遺失或重複處理。
  3. 讓消費速度低於發佈速度,觀察積壓、延遲、容量限制與恢復時間。
  4. 重播指定範圍的事件,確認不會改變已完成結果或重複產生外部效果。
  5. 依照相同鍵值連續送出事件,確認順序、平行處理與擴充結果。
  6. 變更訊息結構,驗證新舊生產者與消費者同時存在時仍能處理。
  7. 由實際維護人員執行告警處理、死信重送與資料修正程序。

POC 應該記錄事件大小與數量、可接受延遲、路由規則、保留期限、重播時間、故障結果及維護工作,再以架構決策紀錄(Architecture Decision Record, ADR)保存選擇與重新評估條件。

建立可以演進的訊息契約

訊息是跨服務的公開契約。欄位名稱、型別、必要性或語意一旦改變,就可能影響不同時間部署的消費者。每一類訊息都要定義擁有者、名稱、結構描述(Schema)與相容策略。

{
    "eventId": "01JEXAMPLE0000000000000000",
    "eventType": "record.updated",
    "occurredAt": "2026-08-22T10:30:00+08:00",
    "producer": "record-service",
    "schemaVersion": 2,
    "subjectId": "R-1042",
    "sequence": 18,
    "operationId": "OP-8f13",
    "data": {
        "status": "active"
    }
}

這個片段只示範事件外層結構,實際欄位要按照追蹤、順序與資料最小化需求選擇:

欄位 用途
eventId 唯一識別事件,供消費者判斷重複處理
eventType 表示已發生的事實,名稱變更時要考慮既有消費者
occurredAt 記錄事件在來源服務發生的時間,不能單獨作為全域順序依據
producer 指出事件契約與內容的負責來源
schemaVersion 讓消費者選擇對應解析方式,不能取代相容設計
subjectId 指出事件所屬對象,也可以作為分割或順序鍵值
sequence 表示同一對象的來源版本,供消費者辨識舊事件或缺漏
operationId 串連觸發操作、發佈、消費與後續處理記錄
data 只放接收責任需要的內容,避免直接公開來源服務的內部保存結構

相容變更通常可以新增選用欄位,並讓消費者忽略不認識的欄位。刪除必要欄位、改變既有欄位語意或型別時,要先讓消費者支援新舊版本,再調整生產者。無法維持相容時,可以建立新的事件類型或目的地,完成轉換後才停止舊版本。每次變更都要用契約測試驗證目前仍會同時運作的版本組合。

事件內容也要遵守資料最小化。消費者只需要識別碼時,不應該附帶完整來源資料。消費者如果必須在來源服務暫時無法使用時完成工作,才評估攜帶必要快照,並定義資料過期與修正方式。

以冪等處理面對重複訊息

可靠傳遞通常會選擇至少一次傳遞(At-Least-Once Delivery)。生產者在結果不明時重新發佈,或消費者在完成處理後、回報進度前停止,都可能讓同一事件再次出現。Apache Kafka 的傳遞語意說明也將發佈與消費分成不同保證範圍,寫入 Kafka 以外的保存機制時仍要協調處理結果與消費位置。

冪等性(Idempotency)表示相同操作執行一次或多次,最終產生相同可觀察結果。消費者可以採用下列方式:

  • 使用 eventId 記錄已處理事件,收到重複事件時直接回覆既有結果。
  • 使用來源版本或 sequence,只套用比目前版本新的事件。
  • 以功能本身的唯一識別碼或條件更新避免重複建立相同結果。
  • 如果消費結果與事件處理紀錄位於同一個交易範圍,將兩者一起提交,避免只完成其中一項。
  • 如果處理會呼叫其他服務,將穩定的冪等鍵傳到下一個責任邊界,並驗證對方如何保存與判斷。

確認訊息或提交消費位置的時機要放在處理成功之後。提前確認可能在後續失敗時遺失工作,延後確認則要接受處理完成後再次收到訊息的可能。工具宣稱的恰好一次處理(Exactly-Once Processing)能力也要確認適用範圍,不能推論它會自動涵蓋訊息代理以外的保存變更或外部效果。

明確限制順序需求

要求所有事件具有全域順序會限制平行處理能力,而且多數功能只需要同一對象的變更依序套用。設計時應該先找出真正需要順序的範圍,再選擇穩定鍵值、分割方式與消費策略。

  • 使用相同 subjectId 或其他擁有者鍵值,讓同一對象的事件進入相同順序範圍。
  • 在訊息內加入來源版本,讓消費者拒絕較舊版本,並在發現版本缺口時暫停或重新取得狀態。
  • 驗證重新排隊、消費者重新分配與重試目的地是否會改變原有順序。
  • 讓彼此無關的對象平行處理,避免為了少數順序需求限制所有事件。

時間戳記不能單獨證明全域先後順序,因為不同執行環境的時間可能有誤差,同一時間也可能發生多項變更。需要正確順序時,應該使用由資料擁有者產生的版本或序號。

區分重試、死信與補償

失敗後立即無限重試,可能讓已失效的相依項目持續承受工作,也可能讓有問題的訊息阻塞後續處理。每類錯誤要先判斷再次執行是否可能成功:

失敗類型 處理方式
短暫連線或相依失敗 使用有次數與總時間上限的退避重試,並加入隨機間隔避免同時重試
處理容量暫時不足 限制消費並行數量,延後重試並監控最舊訊息等待時間
訊息結構或內容無法處理 停止自動重試,保存失敗原因與原始事件識別碼
功能規則拒絕 記錄確定結果,不應該當成暫時性錯誤反覆重試
已產生部分效果 按照流程狀態執行向前修正、反向操作或人工補償

死信佇列(Dead Letter Queue, DLQ)用來隔離超過重試限制或無法處理的訊息。它不能代替錯誤修正,也不能成為無人處理的永久儲存區。RabbitMQ 的死信交換器說明顯示死信仍是一段訊息發佈流程,目的地不可用時也可能失敗。設計時要確認死信原因、保存期限、告警負責人、修正方式、重新送出條件與重送後的冪等結果。

需要跨多個服務完成的流程,可以使用明確狀態記錄每一步是否完成、失敗或等待補償。補償不一定能還原所有外部效果,因此要針對功能語意定義可接受結果,不能把刪除一筆資料當成通用復原方式。

使用交易寄件匣避免雙重寫入缺口

如果服務要在同一項操作中修改自己的資料庫,並且發佈事件通知其他服務,依序執行兩次寫入會留下失敗缺口。資料先提交後,程式可能在事件發佈前停止。事件先發佈後,資料交易則可能失敗,使消費者收到尚未成立的變更。

交易寄件匣模式(Transactional Outbox Pattern)將資料變更與待發佈事件寫入同一個本地交易,再由獨立的訊息轉送程序讀取寄件匣並發佈:

  1. 來源服務開始本地資料庫交易。
  2. 來源服務修改自己擁有的資料,並在寄件匣寫入具有唯一識別碼的事件。
  3. 來源服務提交交易,使資料與事件同時成立,或在失敗時一起取消。
  4. 訊息轉送程序讀取尚未發佈的寄件匣項目,將事件送到訊息代理。
  5. 訊息代理確認接手後,轉送程序記錄發佈結果,並按照保存政策清理寄件匣。
  6. 消費者以冪等方式套用事件,成功後才確認消費進度。

AWS 的交易寄件匣模式說明也將資料更新與寄件匣事件放在同一筆交易,並提醒轉送程序可能重複發佈,因此接收端仍要判斷重複事件。

這個模式解決的是來源資料與「待發佈事件」之間的原子性,不會建立跨服務資料庫交易,也不保證立即一致。轉送程序還要處理查詢頻率、發佈確認、重試、順序、積壓、清理與監控。如果來源保存機制不支援需要的本地交易,就要另外評估變更資料擷取(Change Data Capture, CDC)或能提供相同原子邊界的方式,並用故障測試驗證其保證。

規劃遺留系統與目標系統的切換同步

重構期間可能需要把遺留系統的變更持續同步到目標系統,讓目標系統在正式切換前完成驗證。這段同步是有期限的遷移機制,不應該演變成兩套系統長期互相寫入。

  1. 指定切換前由遺留系統擁有寫入責任,目標系統只接收遷移資料或同步事件。
  2. 建立初始資料快照,記錄快照對應的時間、版本或變更位置。
  3. 從同一個位置擷取快照之後的增量變更,避免快照與即時事件之間出現缺口。
  4. 讓目標系統使用可重複執行的轉換與冪等寫入,並記錄每批資料與事件的處理位置。
  5. 比較資料筆數、必要欄位、關聯、彙總結果與抽樣內容,不能只以訊息已消費判斷同步完成。
  6. 在正式切換前停止或限制遺留系統寫入,追上最後變更並執行驗收。
  7. 切換寫入擁有者後,監控目標系統結果,並按照已演練的停止條件決定繼續或切換回既有系統。

如果切換後仍允許目標系統產生新資料,切換回既有系統前要先定義這些資料如何保留、轉換或限制寫入。缺少這項設計時,單純切換執行版本可能造成切換後資料遺失。雙向同步只有在確實需要且衝突規則可以驗證時才採用,並且要防止同一項變更在兩個方向反覆傳遞。

監控傳遞進度與資料結果

訊息已進入代理只代表傳遞流程開始。監控要同時涵蓋生產、代理、消費與最終資料結果:

  • 記錄發佈成功、發佈失敗、無法路由與發佈確認等待時間。
  • 觀察佇列積壓量、最舊訊息等待時間、主題保留容量與消費者延遲。
  • 記錄消費成功、處理時間、重試次數、死信數量與補償狀態。
  • 使用 eventIdsubjectIdoperationId 串連來源變更、事件發佈、消費與套用結果。
  • 比較來源版本與各衍生副本已套用版本,確認訊息已消費後資料確實符合預期。
  • 為延遲上限、積壓成長、連續失敗、死信與資料差異設定分級告警及處理程序。

監控指標要能回答哪一批資料受影響、最後成功位置在哪裡,以及能否安全重送。只有總訊息數量或程序仍在執行的狀態,無法證明跨服務同步正確。

完成資料同步設計的檢查

  • 是否已指定每項資料的唯一寫入擁有者與其他服務需要的衍生內容?
  • 是否已定義允許延遲、失敗期間結果及最終一致的可檢查條件?
  • 每一步是否按照立即回應需求選擇同步呼叫或非同步訊息?
  • RabbitMQ 或 Apache Kafka 的選擇是否來自路由、保留、重播、順序、處理量與維護限制?
  • 訊息名稱、結構描述、識別碼、版本與新舊消費者相容規則是否明確?
  • 消費者是否能安全處理重複、延遲與順序錯置的訊息?
  • 重試是否具有適用錯誤、次數、退避、總時間與停止條件?
  • 死信與補償流程是否具有負責人、告警、修正、重送及驗證方式?
  • 如果同時修改資料庫並發佈事件,是否以交易寄件匣或經過驗證的替代方式消除雙重寫入缺口?
  • 切換同步是否包含初始快照、增量位置、差異比對、最後同步與切換後資料處理?
  • 監控是否能串連事件,並且確認訊息已被讀取及最終資料結果正確?

重點整理

  • 分散式服務的資料同步要先定義資料擁有者、允許延遲、失敗結果與一致性條件,再選擇通訊方式。
  • 同步 API 呼叫適合必須立即取得結果的步驟,非同步訊息適合允許延遲、分開處理時間或具有多個獨立接收責任的步驟。
  • RabbitMQ 的交換器與佇列適合工作分派及彈性路由,Apache Kafka 的主題與分割區適合保留、重播及讓多個消費者群組獨立處理事件。
  • 訊息契約要定義名稱、結構、識別碼、版本與相容策略,並且只包含接收責任需要的資料。
  • 至少一次傳遞可能產生重複訊息,消費者要以事件識別碼、來源版本或唯一條件建立冪等處理。
  • 順序應該限制在真正需要的對象範圍,並以穩定鍵值與來源版本處理平行執行及順序錯置。
  • 重試、死信與補償分別處理暫時失敗、無法自動處理的訊息及部分完成的流程,每一項都需要明確停止與驗證條件。
  • 交易寄件匣模式讓來源資料與待發佈事件在同一本地交易中成立,轉送與消費端仍要處理重複、積壓及監控。
  • 切換期間應該採用單向、可追蹤且可重複執行的資料同步,完成快照、增量追趕與差異驗證後才轉移寫入責任。

上一篇
[Day 23] 如何讓系統更穩定?需要分散式架構嗎?
系列文
遠古聖遺物改造工程:遺留系統全面重構實務指南24
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言