跳到主要內容

系統設計

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

主題 · 訊息佇列

訊息佇列

生產者把工作丟進佇列就走,消費者按自己的速度慢慢處理:尖峰流量被攤平,一邊掛了另一邊也不受影響。

程式碼路徑
3
60/s
0.5 s
5%
在佇列裡等待的訊息尖峰:每秒 60 則(平時 20 則)
0851700s10s20s30s40s50s60s時間
同步:被拒絕
39.3%
佇列:最長
171
佇列:p99 等待
5.8 s
最後尖峰後清空等候佇列
模擬期間未清空
送達保證:200 則訊息、4 個消費者,5% 的處理中途當機
處理完成
200 / 200
交付次數
216
額外處理次數
8
當機後多久重送
0.7 s

3 個消費者每秒大約處理 30 則,尖峰時每秒進來 60 則。同步呼叫時,1,787 個請求裡有 703 個(39%)遇到所有消費者都在忙而被拒絕。改經過佇列就一則都不會丟:最多 171 則在排隊,最慢的 1% 要等 5.8 s。最後一次尖峰結束後,到模擬結束仍未清空等候佇列。尖峰後仍有新訊息到達;等候佇列清空時,也可能有訊息正在處理。下方:可見性逾時設 0.5 秒時,200 則訊息在首次完成之外,又額外處理了 8 次。消費者只是慢了一點,訊息就可能被交給別人,同一則也可能被處理超過兩次。

假設:請求隨機到達(Poisson),平時每秒 20 則,每 20 秒有 5 秒升到尖峰流量;每個消費者處理一則平均 100 ms(指數分佈)。送達保證那一段,處理一則平均 300 ms,訊息每秒進來 10 則。

亮起來的是這一步執行的程式碼
class VisibilityQueue {
private ready: number[] = []; // waiting, oldest first
private head = 0;
// Received but not yet acknowledged: id -> deadline. Deadlines are
// receive time + a fixed timeout, so the oldest is always first.
private inFlight = new Map<number, number>();
private done = new Set<number>();
constructor(private visibilitySec: number) {}
send(id: number): void {
this.ready.push(id);
}
receive(now: number): number | null {
for (const [id, deadline] of this.inFlight) {
if (deadline > now) break;
this.inFlight.delete(id);
this.ready.push(id);
}
while (this.head < this.ready.length) {
const id = this.ready[this.head++];
if (this.done.has(id)) continue;
this.inFlight.set(id, now + this.visibilitySec);
return id;
}
return null;
}
ack(id: number): void {
this.done.add(id);
this.inFlight.delete(id);
}
}

模型假設與範圍

  • 這是可重現的教學模型;延遲、容量、故障率與工作負載是設定或樣本,不能直接當作正式系統的效能承諾。
  • 工作者、確認、重送與失敗由種子控制;業務副作用和 broker 複寫未實作,確認本身不保證業務恰好執行一次。

什麼時候用

  • 流量不平均:尖峰時進來的工作比消費者處理得快,與其拒絕,不如排隊慢慢消化。
  • 呼叫方不需要等結果:寄信、產生縮圖、寫入分析資料、通知其他服務。
  • 讓兩個服務各自擴縮、各自掛掉:消費者停機時,訊息在佇列裡等它回來。

和其他主題的關係

延伸閱讀
限流器

出現在這些架構裡

時間與空間複雜度(Big O)

操作平均最差
送出O(1)O(1)
取出
均攤 O(1);k 是剛好同時逾時、要放回佇列的訊息數
O(1)O(k)
確認O(1)O(1)

空間:O(n),n 是還沒確認的訊息數;尖峰時佇列就是用這些空間換掉被拒絕的請求

Big O 實測:n 變大時步數怎麼長

數的是:n 則訊息一次湧入、4 個消費者處理

Big On = 100n = 1,000n = 10,000成長倍數:實測(理論)
每次取出的工作量O(1)1.72.02.0×1.1 (×1.0)
清空佇列的時間(秒)O(n)2.524.9259×104 (×100)

取出一則訊息的成本不隨積壓量變多;但消費者數量固定時,清空積壓的時間和積壓量成正比。要更快清空,只能加消費者。

和其他做法比

同步:被拒絕佇列:p99 等待佇列:最長
2 個消費者55%29 s609
3 個消費者39%5.8 s171
4 個消費者28%3.2 s125
6 個消費者15%657 ms31
逾時 0.2 秒額外處理 23 次當機後等 0.6 秒遺失 0 則
逾時 0.5 秒額外處理 8 次當機後等 0.7 秒遺失 0 則
逾時 1 秒額外處理 3 次當機後等 1.3 秒遺失 0 則
逾時 3 秒額外處理 0 次當機後等 3.4 秒遺失 0 則

前四列:同一串尖峰流量(每秒 20 則、尖峰 60 則),同步呼叫和經過佇列各跑一次。後四列:200 則訊息、5% 的處理中途當機,只換可見性逾時。逾時太短,處理慢一點的訊息就被重送、做兩次;太長,當機的訊息要等很久才有人接手。兩種情況都不會遺失訊息——所以消費者必須能安全地處理同一則訊息兩次(冪等)。p99 是第 99 百分位(percentile):前四列只統計每則訊息從到達到消費者開始處理的佇列等待時間,排好後取 99% 位置的值;不含處理本身的耗時,也不是後四列重送嘗試的延遲。

真實世界裡的它

  • Amazon SQS、RabbitMQ、Google Pub/Sub:本節的可見性逾時就是 SQS 的 visibility timeout。
  • Kafka 是可以重播的日誌:消費者自己記住讀到哪裡,訊息讀完也不刪除。
  • Celery、Sidekiq 等背景工作系統,底下都是一個佇列。

取捨與陷阱

  • 至少一次就是可能兩次:消費者要能安全地重做同一則訊息,例如用訊息 ID 去重,或讓操作本身冪等。
  • 佇列只是把等待藏起來:如果平均流量就超過處理能力,佇列會無限變長,最後延遲爆掉或記憶體用完。要監控佇列長度。
  • 一直失敗的訊息會無限重送、卡住後面的人;設重試上限,超過就移到死信佇列(dead-letter queue)。