跳到主要內容

系統設計

把資料結構放大到好幾台機器

主題 · 案例:通知系統

案例:通知系統

一次要送出幾百萬則推播、簡訊和 email:排隊、限流、重試、去除重複,還要尊重每個使用者的通知設定和勿擾時段。

流程
一個活動查偏好到點再送行銷後台業務服務(驗證碼)事件驅動與串流處理通知事件串流使用者偏好與勿擾時段分派 worker延後排程訊息佇列通道佇列(驗證碼優先)冪等與重試冪等鍵(去重)限流器限流與重試發送器APNs/FCM簡訊供應商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,866623023101,000 (1,000)
簡訊1,21367026030 (30)
Email5,1142560950200 (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 On = 100,000n = 1,000,000n = 10,000,000成長倍數:實測(理論)
推播(每秒 1,000)O(n)1001,00010,000×100 (×100)
簡訊(每秒 30)O(n)3,33433,334333,334×100 (×100)

瓶頸是供應商的上限,不是自己的伺服器:多開 worker 沒用,要嘛向供應商申請更高的額度、分給多個帳號或供應商,要嘛接受活動要送好幾個小時,並且把驗證碼的額度保留下來。

和其他做法比

驗證碼等待 p50/p99活動訊息等待 p50/p99收到重複通知的人
優先序+去重0 / 1 s7 / 32402 s0
沒有優先序0 / 23 s7 / 32402 s0
沒有去重0 / 1 s7 / 32402 s361

同一個活動(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 會告訴你哪些權杖已經失效,一直送不但浪費額度,還可能被供應商降級。