昨天做出了最陽春的 Notification System,採用最直接的同步方式:
Charmander -> POST /posts/:id/like -> Create Like -> Create Notification -> Response
這種方式在系統規模小的時候很合理,但功能一多,API 要做的事情就跟著變多。
假設之後除了 Notification,還要加上 Analytics(資料分析),如果所有事情都得在同一個 Request 裡做完才能回應,這支 API 只會越來越慢。
解法是把「不需要讓使用者等待結果」的工作,從同步流程裡抽出去,變成背景作業,這就要靠 Message Queue (MQ)。
可以把 Message Queue 想成一疊待辦清單,紙條上寫著「要做什麼事」,而不是直接叫誰馬上去做。
整體流程是:Producer -> Message Queue -> Consumer
放到我們的場景裡就是:Like Service(Producer)-> MQ(Queue)-> Notification Worker(Consumer)
原本的同步流程:
POST /like
↓
Create Like
↓
Create Notification
↓
Response
改成非同步之後:
POST /like
↓
Create Like
↓
Publish Message
↓
Response

Charmander 送出 Like 之後,只要確認 DB 更新 Like 資料,不需要等 Notification 真的建立完成,Response 幾乎立刻就能回去。剩下的工作交給背景處理:
Message Queue
↓
Notification Worker
↓
Create Notification
丟進 Queue 的 Message 不是一段可執行的程式碼,而是描述「這個工作需要哪些資料」,可以把它想成一張有固定格式的紙條:
實際的內容像這樣:
{
"event": "post_liked",
"post_id": 123,
"actor_id": "Charmander",
"recipient_id": "Pikachu"
}
Notification Worker 收到這張紙條後,就知道「Charmander 對 Post 123 按了讚,要通知 Pikachu」,並依照這些資訊建立通知。
不過事情還是沒那麼簡單,因為解決了一個問題,通常會得到另一組問題(哭啊)
假設 Worker 拿到一張紙條,處理完成、建立了 Notification,卻在跟 Queue 回報「我做完了」之前當機。Queue 不知道這張紙條到底有沒有處理完,只好把它交給另一個 Worker 重做一次,結果同一個事件被處理了兩次。
這牽涉到 Idempotent(幂等性):
同一個操作不管執行一次還是多次,最終結果都應該一樣。
Message Queue 常見的保證是 At-least-once delivery,訊息至少會送達一次但可能重複,所以 Consumer 自己要有辦法應付「被叫做同一件事兩次」的情況。
具體做法可以讓每個 Event 帶一個 event_id,建立 Notification 時把它設成 Unique Constraint。這樣即使同一張紙條被處理兩次,Database 也會擋掉第二次的建立,確保最終只會有一筆 Notification。
如果某張紙條每次交給 Worker 都失敗,也不能讓它無限重試,否則它會一直占用處理資源。
通常會設定一個重試上限,超過之後把這張紙條移到 Dead Letter Queue,一個專門收留「暫時處理不了」的紙條的地方,之後工程師可以慢慢查原因、修好問題、決定要不要重新丟回 Queue。
假設 Producer 每秒產生 1000 張紙條,Worker 每秒卻只能處理 500 張,清單只會越疊越高。
這時候需要盯著兩個指標:
判斷原則很簡單:使用者是否需要在這次 Request 裡立刻拿到結果?
適合放進 Queue、稍後處理:
不適合非同步處理:
現在可以再把格局拉高一層。post_liked 這個事件,其實不是只有 Notification Worker 需要知道,Analytics、Counter 也可能都對它感興趣:

一個事件觸發多個不同的背景工作,這就是 Event-driven Architecture 的基本概念。Like Service 完全不需要知道「到底有哪些服務要處理 post_liked」,各個服務自己決定要不要訂閱這個事件,這讓系統可以很輕鬆地加入新的 Consumer,而不必去修改原本的 Producer。
透過 Message Queue,我們把 Notification 從主流程裡拆了出去:
使用者不再需要為了一則通知而多等一段時間,而系統也因為 Producer 與 Consumer 解耦,變得更容易個別擴展、也更能承受突發流量。
到目前為止,PokeThreads 已經具備了 Load Balancer、Stateless Server、Session Store、Cache、Read Replica、CDN、Media Pipeline 與 Message Queue,完整架構長這樣:

是時候回頭想想,這些元件之間該怎麼有效率地互相溝通了。