訊息佇列
生產者把工作丟進佇列就走,消費者按自己的速度慢慢處理:尖峰流量被攤平,一邊掛了另一邊也不受影響。
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 O | n = 100 | n = 1,000 | n = 10,000 | 成長倍數:實測(理論) | |
|---|---|---|---|---|---|
| 每次取出的工作量 | O(1) | 1.7 | 2.0 | 2.0 | ×1.1 (×1.0) |
| 清空佇列的時間(秒) | O(n) | 2.5 | 24.9 | 259 | ×104 (×100) |
取出一則訊息的成本不隨積壓量變多;但消費者數量固定時,清空積壓的時間和積壓量成正比。要更快清空,只能加消費者。
和其他做法比
| 同步:被拒絕 | 佇列:p99 等待 | 佇列:最長 | |
|---|---|---|---|
| 2 個消費者 | 55% | 29 s | 609 |
| 3 個消費者 | 39% | 5.8 s | 171 |
| 4 個消費者 | 28% | 3.2 s | 125 |
| 6 個消費者 | 15% | 657 ms | 31 |
| 逾時 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)。