案例:通知系統
一次要送出幾百萬則推播、簡訊和 email:排隊、限流、重試、去除重複,還要尊重每個使用者的通知設定和勿擾時段。
流程
通知系統把「發一個活動」變成幾萬次對外部供應商的呼叫,每一家都有自己的每秒上限,也會偶爾失敗。中間每一段都隔著佇列,各段用自己的速度跑。選一個流程,看勿擾時段、驗證碼優先和去重分別在哪一段發生。
亮起來的是這一步執行的程式碼
type User = { id: number; utcOffset: number; channel: "push" | "sms" | "email" | "none" };type Message = { key: string; urgent: boolean; attempt: number };type Result = "sent" | "duplicate" | "retry" | "failed" | "deferred"; const QUIET_START = 22, QUIET_END = 8, MAX_ATTEMPTS = 4; function localHour(user: User, nowUtcHours: number): number { return (((Math.floor(nowUtcHours) + user.utcOffset) % 24) + 24) % 24;} function plan(user: User, nowUtcHours: number, urgent: boolean): "skip" | "send" | "defer" { if (user.channel === "none") return "skip"; if (urgent) return "send"; const hour = localHour(user, nowUtcHours); return hour >= QUIET_START || hour < QUIET_END ? "defer" : "send";} function backoffSeconds(attempt: number): number { return 2 ** (attempt - 1);} function nextAllowedUtc(user: User, nowUtcHours: number): number { const local = ((nowUtcHours + user.utcOffset) % 24 + 24) % 24; return local >= 8 && local < 22 ? nowUtcHours : nowUtcHours + ((8 - local + 24) % 24);} class Dispatcher { private high: Message[] = []; private low: Message[] = []; private delivered = new Set<string>(); private users = new Map<string, User>(); private later: Array<{ at: number; message: Message }> = []; // Store deferred work; the caller advances UTC time in sendSecond(). schedule(user: User, message: Message, nowUtc: number): "skip" | "send" | "defer" { const decision = plan(user, nowUtc, message.urgent); if (decision === "skip") return decision; this.users.set(message.key, user); if (decision === "defer") this.later.push({ at: nextAllowedUtc(user, nowUtc), message }); else this.enqueue(message); return decision; } enqueue(message: Message): void { (message.urgent ? this.high : this.low).push(message); } // One second of sending to a provider that allows `limit` calls. sendSecond(limit: number, provider: (m: Message) => boolean, nowUtc = 0): Result[] { const due = this.later.filter(x => x.at <= nowUtc); this.later = this.later.filter(x => x.at > nowUtc); for (const x of due) this.enqueue(x.message); const results: Result[] = []; let calls = 0; while (calls < limit) { const message = this.high.shift() ?? this.low.shift(); if (!message) break; if (this.delivered.has(message.key)) { results.push("duplicate"); continue; } const user = this.users.get(message.key); if (user && this.scheduleDecision(user, message, nowUtc)) { results.push("deferred"); continue; } calls++; if (!provider(message)) { results.push(message.attempt < MAX_ATTEMPTS ? "retry" : "failed"); continue; } this.delivered.add(message.key); results.push("sent"); } return results; } private scheduleDecision(user: User, message: Message, nowUtc: number): boolean { if (plan(user, nowUtc, message.urgent) !== "defer") return false; this.later.push({ at: nextAllowedUtc(user, nowUtc), message }); return true; }}04:00
驗證碼優先
冪等鍵去重
5%
時間以 UTC(Coordinated Universal Time)計。一個活動發給 20,000 人(60% 在 UTC+8,其餘在東京、倫敦、紐約、洛杉磯;用標準時間,不算夏令時間),另外前 120 秒每秒有 2 組簡訊驗證碼。勿擾時段是當地 22:00–08:00。供應商每秒上限:推播 1,000、Email 200、簡訊 30。失敗後隔 1、2、4 秒重試;2% 送出的訊息會再從佇列出來一次。現在:台北 12:00、倫敦 04:00、紐約 23:00。
延後(勿擾時段)
3,627
模擬結束時仍待送
0
勿擾時段內送達
0
驗證碼等待 p99
1 s
已送達活動訊息等待 p99
32402 s
| 送達 | 重試 | 放棄 | 擋下的重複 | 重複送達 | 每秒最多(上限) | |
|---|---|---|---|---|---|---|
| 推播 | 11,866 | 623 | 0 | 231 | 0 | 1,000 (1,000) |
| 簡訊 | 1,213 | 67 | 0 | 26 | 0 | 30 (30) |
| 5,114 | 256 | 0 | 95 | 0 | 200 (200) |
3,627 則行銷訊息因勿擾時段延後,沒有任何一則在勿擾時段送達。驗證碼插隊到簡訊佇列最前面:99% 在 1 秒內送達。被重複投遞的訊息都被冪等鍵擋下,沒有人收到兩次同樣的通知。
亮起來的是這一步執行的程式碼
type User = { id: number; utcOffset: number; channel: "push" | "sms" | "email" | "none" };type Message = { key: string; urgent: boolean; attempt: number };type Result = "sent" | "duplicate" | "retry" | "failed" | "deferred"; const QUIET_START = 22, QUIET_END = 8, MAX_ATTEMPTS = 4; function localHour(user: User, nowUtcHours: number): number { return (((Math.floor(nowUtcHours) + user.utcOffset) % 24) + 24) % 24;} function plan(user: User, nowUtcHours: number, urgent: boolean): "skip" | "send" | "defer" { if (user.channel === "none") return "skip"; if (urgent) return "send"; const hour = localHour(user, nowUtcHours); return hour >= QUIET_START || hour < QUIET_END ? "defer" : "send";} function backoffSeconds(attempt: number): number { return 2 ** (attempt - 1);} function nextAllowedUtc(user: User, nowUtcHours: number): number { const local = ((nowUtcHours + user.utcOffset) % 24 + 24) % 24; return local >= 8 && local < 22 ? nowUtcHours : nowUtcHours + ((8 - local + 24) % 24);} class Dispatcher { private high: Message[] = []; private low: Message[] = []; private delivered = new Set<string>(); private users = new Map<string, User>(); private later: Array<{ at: number; message: Message }> = []; // Store deferred work; the caller advances UTC time in sendSecond(). schedule(user: User, message: Message, nowUtc: number): "skip" | "send" | "defer" { const decision = plan(user, nowUtc, message.urgent); if (decision === "skip") return decision; this.users.set(message.key, user); if (decision === "defer") this.later.push({ at: nextAllowedUtc(user, nowUtc), message }); else this.enqueue(message); return decision; } enqueue(message: Message): void { (message.urgent ? this.high : this.low).push(message); } // One second of sending to a provider that allows `limit` calls. sendSecond(limit: number, provider: (m: Message) => boolean, nowUtc = 0): Result[] { const due = this.later.filter(x => x.at <= nowUtc); this.later = this.later.filter(x => x.at > nowUtc); for (const x of due) this.enqueue(x.message); const results: Result[] = []; let calls = 0; while (calls < limit) { const message = this.high.shift() ?? this.low.shift(); if (!message) break; if (this.delivered.has(message.key)) { results.push("duplicate"); continue; } const user = this.users.get(message.key); if (user && this.scheduleDecision(user, message, nowUtc)) { results.push("deferred"); continue; } calls++; if (!provider(message)) { results.push(message.attempt < MAX_ATTEMPTS ? "retry" : "failed"); continue; } this.delivered.add(message.key); results.push("sent"); } return results; } private scheduleDecision(user: User, message: Message, nowUtc: number): boolean { if (plan(user, nowUtc, message.urgent) !== "defer") return false; this.later.push({ at: nextAllowedUtc(user, nowUtc), message }); return true; }}模型假設與範圍
- 這是可重現的教學模型;延遲、容量、故障率與工作負載是設定或樣本,不能直接當作正式系統的效能承諾。
- UTC offset 固定,不含日光節約;延後訊息排到下一個當地 08:00,送出前再檢查。預設跑 24 小時,期限內未完成者列待送;延後是不同訊息數,p99 只算已送達訊息,從原始建立時間開始。展示 dispatcher 的 retry 結果由呼叫端安排重試。
什麼時候用
- 每個通道一條佇列、一組發送器:簡訊供應商慢或掛了,推播和 Email 照常送。
- 緊急訊息走優先佇列:OTP(One-Time Password)、資安警示排在行銷訊息前面,必要時保留一部分供應商額度給它們。
- 勿擾時段用收件人的時區算,而且送出前再檢查一次;延後的訊息交給排程,不要在 worker 裡睡著等。
- 送出前用冪等鍵去重(例如活動編號+使用者):佇列是至少一次送達,重試和重新投遞都會造成重複。
和其他主題的關係
時間與空間複雜度(Big O)
| 操作 | 平均 | 最差 |
|---|---|---|
| 展開活動 n 是收件人數;每人查一次偏好 | O(n) | O(n) |
| 決定送/延後/略過 | O(1) | O(1) |
| 進出優先佇列 | O(1) | O(1) |
| 冪等鍵檢查 雜湊表;最壞情況是碰撞全落在同一格 | O(1) | O(n) |
空間:O(n),每則送出的訊息留一把冪等鍵,過期後刪除
Big O 實測:n 變大時步數怎麼長
數的是:發完一個活動要的秒數(n 是收件人數,只用一個通道)
| Big O | n = 100,000 | n = 1,000,000 | n = 10,000,000 | 成長倍數:實測(理論) | |
|---|---|---|---|---|---|
| 推播(每秒 1,000) | O(n) | 100 | 1,000 | 10,000 | ×100 (×100) |
| 簡訊(每秒 30) | O(n) | 3,334 | 33,334 | 333,334 | ×100 (×100) |
瓶頸是供應商的上限,不是自己的伺服器:多開 worker 沒用,要嘛向供應商申請更高的額度、分給多個帳號或供應商,要嘛接受活動要送好幾個小時,並且把驗證碼的額度保留下來。
和其他做法比
| 驗證碼等待 p50/p99 | 活動訊息等待 p50/p99 | 收到重複通知的人 | |
|---|---|---|---|
| 優先序+去重 | 0 / 1 s | 7 / 32402 s | 0 |
| 沒有優先序 | 0 / 23 s | 7 / 32402 s | 0 |
| 沒有去重 | 0 / 1 s | 7 / 32402 s | 361 |
同一個活動(20,000 人,UTC 04:00 開始)、同一組故障,每次只關掉一項。優先序只決定誰先送,不會讓整個活動變慢;去重的代價是每則訊息查一次鍵。
真實世界裡的它
- iOS 的推播要經過 APNs(Apple Push Notification service),Android 經過 FCM(Firebase Cloud Messaging);伺服器不能直接連到手機。
- Twilio、Amazon SNS(Simple Notification Service)、SendGrid 這類服務都有每秒或每日的發送上限,超過會回 429 要你慢下來。
- 電商的大型促銷常把推播分批在幾十分鐘內送完,避免所有人同時點開把網站打掛。
取捨與陷阱
- 用伺服器時區算勿擾時段:台北中午送出的活動,在紐約是半夜。
- 失敗就立刻重試:供應商正在限流時,立刻重試只會讓它更忙;要指數退避,加上隨機抖動,並設上限次數。
- 冪等鍵用隨機產生的訊息編號:每次重試編號都不同,等於沒有去重;鍵要從業務意義來(活動+使用者、驗證碼請求編號)。
- 沒有處理退訂和失效的裝置權杖:APNs/FCM 會告訴你哪些權杖已經失效,一直送不但浪費額度,還可能被供應商降級。