假設一筆事件已經處理完,Consumer 卻在 ack 前中斷。重啟之後,同一筆又送來了,訂單或庫存會不會再做一次?我想先把這個情境想清楚,再談事件流程怎麼實作。
採至少一次投遞,就要準備接到重複事件。處理完還沒 ack、連線中斷、重啟後重新消費,都可能遇到。所以,我會把「同一件事再來一次,也不會多做一次副作用」當成流程的基本條件。
另一邊也要照顧:資料已經更新,事件卻沒有送出去。下游不知道有這件事,就不會繼續處理。送出端與接收端,都需要各自準備。
訂閱者看到名字,要能知道收到的是什麼結果,再決定自己接下來要做什麼。所以事件命名,我會用已經發生的事來描述,例如 OrderCreated。

圖 Day 19-1:冪等事件處理。
看圖時,可以順著 Consumer 的順序走一次:先確認 eventId 有沒有處理過。已經處理過就只回 ack、不重做業務更新;還沒處理過,才把業務更新與 inbox 紀錄放進同一筆交易,成功後才 ack。交易成功但 ack 沒送到,重送時就會走上面那條分支;交易失敗,兩邊也不會只留下一半。
Schema 異動也要維持相容,Consumer 需依契約處理欄位差異。新增欄位應為選填,移除或改變既有欄位語意則需要新的 schemaVersion,並確認所有訂閱者都已可處理,才能停止發佈舊版本。
實作上,我會用事件或單據的唯一鍵搭配資料庫交易來處理冪等。原因是兩個 Consumer 或兩個請求有可能在同一時間都查到「尚未處理」,如果只靠程式先查一次,很容易在併發時一起通過。因此,前面的查詢比較像是快速檢查,真正守住最後一道防線的,還是資料庫的唯一約束與交易。
在 relay 這條線上,我使用的唯一鍵不是事件 ID,而是由 source_system、doc_type、source_doc_no 三個欄位組成的業務鍵。原因是對 relay 來說,我真正想避免的是「同一張來源單據被收進來兩次」,而不是單純避免同一個 event ID 重複。
兩層檢查的部分,簡化後大概是這樣:
// txn_relay_log 上有
// UNIQUE (source_system, doc_type, source_doc_no)
// 真正處理併發重複的最後一道防線
public Optional<String> findExistingRelay(
String sourceSystem, String docType, String sourceDocNo) {
String key = "relay:" + sourceSystem + ":" + docType + ":" + sourceDocNo;
try {
// 第一層:Redis 快速檢查,TTL 24 小時
String cached = redis.opsForValue().get(key);
if (cached != null) {
return Optional.of(cached);
}
} catch (Exception e) {
// Redis 異常就繼續查資料庫,不因快取故障拒收資料
log.warn("Redis idempotency check failed (fail-open): {}", e.getMessage());
}
// 第二層:查資料庫
return relayLogRepository
.findBySourceSystemAndDocTypeAndSourceDocNo(sourceSystem, docType, sourceDocNo)
.map(RelayLog::getRelayId);
}
呼叫端收單時,會先透過 findExistingRelay 檢查這筆單據是否已經收過。查到既有資料時,兩種來源的處理不一樣:
| 來源 | 查到已經收過 | 為什麼這樣處理 |
|---|---|---|
| API 呼叫端 | 回 HTTP 409,附上原本的 relayId |
讓呼叫端知道這筆已經收過、不必重送 |
| 檔案批次匯入 | 跳過那一筆,繼續處理後面的資料 | 沒有一個即時的呼叫端等著收 HTTP 回應 |
但這裡還有一個問題:前面的 Redis 與資料庫查詢,都不能真正避免併發。
例如兩個請求幾乎同時進來,請求 A 與請求 B 都可能先查 Redis,結果都查不到;接著再查資料庫,也都還查不到。這時候兩邊都會認為自己可以寫入。
真正決定結果的是資料庫上的:
UNIQUE (source_system, doc_type, source_doc_no)
其中一筆會先成功寫入,另一筆在寫入時撞上唯一約束,整筆交易被退回,什麼都沒有留下。
程式接住這個例外之後,我不會直接把它當成一般錯誤回出去,而是再查一次資料庫。如果這時候可以找到相同業務鍵的 relay 紀錄,就代表剛才確實是併發造成的重複寫入,因此取得第一筆的 relayId 後,一樣回 409 給呼叫端。
如果補查之後仍然找不到資料,才把原本的例外繼續往外丟。原因是資料庫丟出的約束例外不一定都是重複單據,例如欄位空值、外鍵不符合或其他資料完整性問題,也可能走到相同的例外處理。
Producer 這一側,哪些事要放在同一筆交易裡,有兩個方向相反的決定。一是 transactional outbox:我會把 relay 紀錄與待發佈的 outbox event 一起寫進去,避免業務資料已經成功寫進資料庫,但事件還沒來得及送出去,系統就中斷。二是這組業務鍵,反而要等交易成功之後才寫回 Redis。因為 Redis 不在資料庫交易裡,如果在交易中途就寫進去,交易後來被退回時那個鍵會留著,下一次同一張單就會被誤判成重複。
不過 transactional outbox 解決的是「資料與待發佈事件的一致性」,不代表事件只會送一次。例如 relay 已經把事件成功送到 Kafka 或 RabbitMQ,但還沒來得及把 outbox 狀態更新成 SENT 就中斷,服務重新啟動之後,這筆事件就有可能再次被發佈。
因此後面的 Consumer 還是要自己處理重複。一般作法是在 Consumer 端搭配 inbox 或自己的唯一鍵機制,記錄哪些事件已經處理過;收到事件時先確認是否已經執行,如果已經處理過就直接略過。
| 取捨 | 這樣選的理由 | 何時要重新評估 |
|---|---|---|
| 以資料庫唯一鍵擋重複 | 避免應用層先查再寫的 race condition | 無 |
| 業務更新與 inbox 同交易 | 「已處理」與業務結果必須一致 | 無 |
| producer 端採 outbox | 避免資料已改但事件未送 | 業務可接受事件遺失時 |
多了 inbox 與 outbox,就多了資料表與監控。尤其 outbox 堆積時,可能是下游還沒收到,也可能是已經送出但來不及標記,兩種都要有人看,所以我會替它安排獨立告警與處理方式。
這個流程,我會主動安排中斷來練習,確認下面幾個情境:
| 驗證項目 | 做法 | 想確認什麼 |
|---|---|---|
| 重複投遞 | 在 DB 更新完成、尚未 ack 時終止程序,再重送同一事件 | 是否產生第二次業務副作用? |
| 並行重複 | 同時送進兩筆相同業務鍵的單據 | 唯一鍵是否真的擋下? |
| outbox 中斷 | 在 relay 發佈後、標記前中斷 | 事件是否重複送出且被 consumer 擋下? |
| schema 演進 | 以新舊版本事件同時投遞 | Consumer 是否依契約容錯? |
後面持續看幾個指標:重複事件、冪等命中率(收到的單據裡,有多少是查到已經收過而直接擋下的)、完成時間、outbox 堆積量與 DLQ。如果冪等命中率一直是零,也可以再跑一次重複投遞測試,確認去重機制確實有接上。
今天想完成的,是讓流程遇到重複事件時,仍然能把同一件事處理正確。下一篇,把這份準備延伸到外部夥伴,看看 EDI 與 B2B 還需要照顧哪些細節。