冪能消費者模式

設計訊息消費者,讓同一則訊息即使被處理多次,其效果也與只處理一次相同。 保證至少一次送達的訊息系統,可以多次傳送同一訊息。 防重複的韌性確保重新處理訊息不會產生重複紀錄、不會對客戶重複收費,或產生其他不良影響。

內容和問題

分散式應用程式通常透過訊息中介交換工作,而非直接同步呼叫。 大多數經紀商,包括 Azure 服務匯流排、Azure 事件中樞、Apache Kafka 和 RabbitMQ,至少提供一次交付服務。 此保證確保即使發生故障,訊息仍能送達消費者,同時也意味著經紀人可以多次傳送同一訊息。

在分散式系統中保證精確 一次 的傳送是不切實際的。 即使是使用恰好一次語義的訊息代理程式,也只能保證其直接控制的操作,例如將訊息傳遞給消費者,或將資料寫回代理程式。 他們無法控制消費者在外部系統中實施的副作用。 持久的解決方案不是消除重複配送,而是讓消費者正確處理。 當你將至少一次交付與忽略重複郵件的用戶結合時,你就能實現 一次 有效的處理。

重複可能來自多種來源:

  • 製作人會重試。 製作者發送訊息後,因暫時性網路故障或逾時而未收到確認訊息,然後再次發送訊息。 經紀人現在持有兩份副本,儘管第一次發送成功。

  • 遺失確認後的重新送達。 消費者接收並處理訊息,但因用戶當機、鎖定到期或確認遺失而未能確認。 經紀人會假設訊息未被處理,然後再次送達。

  • 消費者在處理過程中失誤。 消費者完成資料庫寫入後,在確認訊息前就當機了。 另一個實例會接收訊息並重複寫入。

Solution

透過讓消費者保留其已成功處理之訊息的紀錄,並根據在重新傳遞後仍保持不變的穩定識別碼跳過任何已處理過的訊息,來建立冪等消費者。 消費者會檢查持久識別碼儲存庫,判斷是否已處理該識別碼,然後處理該訊息或將其視為重複丟棄。

核心流包含以下步驟:

  1. 閱讀訊息並提取其去重金鑰。
  2. 檢查去重複儲存區中是否有該金鑰。
  3. 如果該金鑰已經存在於儲存庫中,請將訊息視為重複。 確認訊息並停止處理,並可選擇返回先前錄製的結果。
  4. 如果金鑰還不存在於儲存庫中,就先處理訊息並以單一原子操作記錄金鑰,然後確認訊息。

以下章節提供讓消費者具備冪等性的指導:

選擇穩定的重複去重金鑰

去重金鑰必須在每次重送中唯一且一致地識別邏輯訊息。 使用生產者指派的訊息識別碼或企業層級的冪性鍵來識別特定的邏輯操作,而非多個訊息可能攜帶的共享關聯上下文。

例如,將 服務匯流排 MessageId 屬性設定為一個能唯一識別邏輯訊息的值。 不要使用 CorrelationId 作為鍵,因為它指的是一組相關訊息,例如某個請求及其回覆。 對於遵循 CloudEvents 規範的事件, source 與 id 屬性組合能唯一識別事件,並在重送過程中保持穩定。

不要依據代理程式在重新傳遞時會重新產生的傳輸層級識別碼,或依據根據傳遞嘗試次數衍生出的值。 這些數值會在不同配送間變動,從而避免重複偵測。 同時也避免從易失性欄位(如接收時間戳)中推導金鑰。

當多個獨立消費者處理同一通道時,例如在發佈-訂閱設計中,每個消費者會收到自己的訊息副本,並需獨立追蹤處理完成度。 如果消費者共用一個去重儲存庫,則將記錄鍵化為消費者身份與訊息身份的複合狀態。 僅以訊息身份為鍵的儲存裝置,會讓第一個消費者抑制所有其他訊息的處理。

決定處理後的金鑰存放在哪裡

以下儲存選項是處理中金鑰常見的:

  • 專用去重表。 消費者會維護一個獨立的資料表,有時稱為 收件匣,其中每個已處理的金鑰各佔一列。 此方法將重複刪除問題與業務資料分開,且若多種訊息類型共享相同機制,效果良好。

  • 商業實體本身。 消費者將金鑰儲存在訊息所建立或更新的記錄中。 此方法可避免使用另外的資料表,但會將去重機制與業務資料類型耦合在一起。

以原子方式提交已處理的金鑰及副作用

檢查再處理流程有一個失敗視窗。 如果消費者處理完訊息後,再於另一個步驟記錄該金鑰,那麼如果在這兩個操作之間發生當機,副作用已經生效,但金鑰尚未被記錄,因此消費者會在下一次再次投遞時重新處理該訊息。

透過在同一筆交易中寫入去重標記和業務副作用,避免出現此故障窗口。 如果消費者要嘛將兩個操作一併提交,要嘛完全不提交,那麼在重新投遞時,它不是會找到標記並略過該訊息,就是會因為該交易未完成而找不到標記,並可安全地重新處理該訊息。 此交易變體稱為 收件匣 模式,是生產者 交易寄件箱模式的消費者端伴侶。

防止同時出現的重複

在至少一次傳遞且存在並行競爭消費者的情況下,兩個實例可能會同時收到同一則訊息的副本。 兩個實例都能在任一實例提交前通過存在性檢查,因此單靠檢查並不能防止重複處理。

透過以下步驟,在資料儲存區而非應用邏輯中強制執行正確性:

  • 對去重金鑰使用唯一性限制,使得兩個交易可以嘗試插入一個金鑰,但只有一個能成功。 另一筆交易未達成限制,將訊息視為重複。 這種做法使資料庫成為衝突的唯一仲裁者。

  • 避免快取中先檢查再設定的競態。 一個檢查金鑰並將它設為兩個獨立操作的模式,會有一個視窗允許同時重試來認領該金鑰。 使用原子條件寫入,例如發生衝突時會失敗的插入,或僅在不存在時才設定的操作,使取得該鍵成為單一原子步驟。

處理無法加入交易的副作用

有些處理程序無法參與消費端的資料庫交易作業,例如呼叫第三方 API 或寫入外部儲存體。 這些流程可採用以下兩階段方法:

  1. 先將該金鑰記錄為 進行中 狀態,然後執行外部操作。
  2. 更新 紀錄以完成 並儲存結果。

重新交付時, 完成 的紀錄會告訴消費者跳過重複通話。 進行中的記錄表示先前的嘗試可能已部分完成,或正由其他使用者處理。 消費者應先調和過時的紀錄,或將未解決案件轉介給處理,再確認重新送達。

問題和考慮

在決定如何實施此模式時,請考慮以下幾點:

  • 偏好自然冪등 運算。 有些操作本質上是冪分的,不需要重複去帳。 以商業識別碼為鍵的 upsert、設定絕對值而非遞增的寫入,或是對資源識別碼的 HTTP PUT 執行,無論執行一次或多次,產生相同的結果。

    有時你可以透過由事件攜帶狀態的傳遞,讓某項操作自然地具備冪等性。 訊息攜帶產生的絕對狀態,例如訂單的新狀態,因此消費者將其應用為上溢而非相對變更。

    小提示

    如果可能,應設計成具備天然冪等性,並且僅對無法設計成天然冪等的操作使用去重複技術。

  • 使用訊息框架代替配置重複刪除。 正確實作去重儲存、提交和清理很容易出錯。 基於訊息的框架將此模式作為內建功能提供。

    例如, NServiceBus 透過訊息識別碼來解重訊息,並提供可配置的重複資料保留與清理功能。 MassTransit 消費者寄件匣透過訊息識別碼追蹤收到的訊息,提供精確一次的消費者行為。

  • 管理重複資料刪除紀錄的生命週期。 除非你將重複資料刪除記錄設為到期,否則這些記錄會持續累積。 至少在經紀人還能重傳原始訊息時,保留每筆紀錄。 此視窗大小取決於經紀人的最大投遞嘗試次數、鎖定或可見逾時,以及訊息存活時間。

    對重複刪除紀錄設定超過此期限的存活時間,這樣延遲的重新送達仍能找到標記。 太早刪除紀錄會重新開啟重複的視窗。 將操作員從死信佇列重新提交的訊息納入考量,因為這類重新提交可能會在正常的重新傳遞時限結束很久之後才發生。

  • 不要用訊息代理程式的去重機制取代消費者端的冪等邏輯。 有些平台會在傳輸層過濾重複項目。 例如,服務匯流排 重複偵測會捨棄在已設定的時間範圍內重複某個 MessageId 的訊息,從而抑制生產者重複傳送的重試。

    此功能在發送端運作,且設於有界視窗內,因此不會阻止消費者在重投後重複處理同一則訊息。 你仍然需要消費者端的冪等邏輯。 使用平台功能來減少重複量,而非取代冪等消費者模式。

  • 將訊息排序納入考量。 去除重複會移除重複項目,但不保證順序。 若消費者依賴處理順序,則將此模式與排序機制(如 服務匯流排 訊息會話)結合,或包含序列或版本資料,讓消費者能拒絕過時訊息。

  • 可觀察性儀器。 在結構化日誌中發送去重金鑰與相關識別碼,並追蹤偵測到的重複指標。 重複率上升可能表示生產者設定不當、確認或鎖定視窗過小,或消費者運作異常。 使用 分散式追蹤與關聯 來跨服務追蹤訊息。

  • 將冪等性傳遞至下游呼叫。 讓訊息消費者具備冪等性,並不能保護它所呼叫的服務。 當消費者在處理過程中呼叫下游服務時,傳播冪등 性鍵,讓每個服務層級都能自行解重工作。

使用此模式的時機

當下列情況時,請使用此模式:

  • 你從提供至少一次送達服務的經紀商接收訊息,這是大多數經紀商的預設方式。

  • 重新處理訊息可能導致錯誤結果,例如重複的財務交易、重複的資源建立或重複通知。

  • 多個競爭的消費者同時處理同一通道,這使得同時重複傳送的可能性增加。

在下列情況下,此模式可能不適用:

  • 消費者執行的操作本身就具有冪元性,因此再處理是無害的,而重複去碼記帳只會增加成本卻沒有好處。

  • 工作負載可以容忍偶爾發生的重複處理所帶來的影響,而且去重儲存體的成本高於單次重複所造成的影響。

訊息之外的冪能處理

此模式將冪等性套用於訊息取用者,但冪等處理是一項更廣泛的可靠性原則,對任何針對同一項工作重複執行的作業都有幫助。 此原則包括會重新處理重播資料的擷取、轉換、載入(ETL)作業、從檢查點繼續執行的串流處理、會重疊執行或重新啟動的排程作業,以及會接收重複傳遞內容的 webhook 或 HTTP 端點。

每個案例都適用相同的核心技術。

  1. 使用穩定鍵來識別工作單位。
  2. 記錄你處理的內容。
  3. 跳過或吸收重複的流程,避免重複工作改變結果。

本文中的機制,如穩定金鑰、原子標記和唯一性限制,即使沒有訊息中介,也會轉移到這些上下文中。

工作負載設計

評估如何在工作負載設計中使用冪能消費者模式,以達成Azure Well-Architected框架支柱所涵蓋的目標與原則。 下表提供此模式如何支援每個要素目標的指引。

支柱 此模式如何支援支柱目標
可靠性 設計決策有助於使工作負載具有韌性,並確保在故障發生後能復原到正常運作的狀態。 此模式讓工作負載至少能使用一次傳遞並安全重試且不會損壞資料,將重複交付從正確性風險轉變為可容忍的狀態。

- RE:07 自我保護
- 暫時性錯誤

如果此模式在一個支柱內部引入取捨,請將它們與其他支柱的目標進行考量。

Example

以下範例說明一個冪等取用者,該取用者會處理來自 服務匯流排 的訂單,並將狀態保存於 Azure Cosmos DB for NoSQL 中。

  1. 生產者將 服務匯流排 MessageId 設定為業務層級訂單識別碼。
  2. 消費者會以 PeekLock 模式接收訊息,這會使訊息在消費者未於鎖定到期前完成處理時可再次傳遞。
  3. 消費者的 Azure Cosmos DB 容器會對/orderId訂單識別碼進行分割,並將文件id設定為相同的訂單識別碼,因此每個訂單的副本都會解析到相同的邏輯分割區,而訂單id本身就作為重複刪除標記。

消費者會採取以下步驟來處理每則訊息:

  1. 閱讀訊息並用它 MessageId 作為去重金鑰。
  2. 嘗試建立訂單文件,且將 id 和 分割鍵都設定為訂單識別碼。
  3. 若建立成功,則完成訊息,讓 服務匯流排 將其從佇列中移除。
  4. 如果建立因 HTTP 409(衝突)狀態碼而失敗,原因是已存在具有該 id 的文件,請讀取現有文件,並將其與目前訊息進行比較。
  5. 如果儲存的請求雜湊值或不可變的業務欄位相符,則將訊息視為重複,完成該訊息並跳過後續處理。
  6. 如果儲存的請求雜湊值或不可變的商業欄位不匹配,請將訊息送入死信佇列並發出警示,而不是默然丟棄訊息。 製作人可能重複使用該識別碼用於不同內容,或訊息細節自訂單首次處理以來有所變動。
  7. 若因暫時性原因處理失敗,應放棄訊息,讓 服務匯流排 重新傳送,或讓鎖定過期,讓其他使用者接收。

Create 操作是原子性的,因此同時擔任去重檢查和寫入操作。 兩個收到相同訊息副本的消費者,不能同時建立訂單。 一次創建成功,另一次嘗試則返回衝突並安全丟棄其複製品。

當處理程序必須寫入一份以上的文件時,請使用 交易式批次,其中同時包含去重複金鑰與業務文件,且兩者皆位於相同的分割區索引鍵下。 由於交易批次是在單一邏輯分割區內運作,你選擇一個分割鍵,讓同一訊息的所有文件共享。 批次會一次提交所有文件,否則全部都不提交,因此即使在處理與確認之間發生當機,也不會讓去重標記與業務資料失去同步。若某個批次嘗試建立已存在的文件,則會回傳 409(Conflict)狀態碼,藉此識別重複項。

為了讓這個冪等消費者也能因應重複傳送重試,請在佇列上啟用重複偵測。 在標準或高級隊列中,重複偵測會抑制其歷史視窗內的重複發送。 冪等消費者仍然會處理任何超出該時間窗範圍,或因重新傳遞而產生的重複訊息。

後續步驟