Low-level design 的並發發生在 一個 process 裡。Threads 共享同一份 heap。要看清的是:兩條 thread 重疊時什麼會壞掉,以及在發明第四種機制之前,該伸手去拿哪一種 primitive。
這篇筆記依 Hello Interview。書面對照是 concurrency intro。JavaScript 與 TypeScript 跑在 main thread 上是例外 —— user code 並不跨 thread 共享記憶體;那一篇是 深入理解 JavaScript Event Loop。這篇筆記講的是分類。
Pattern Map
| 類型 | 壞在哪裡 | Primitive | 出現在 |
|---|---|---|---|
| Correctness | Shared state 被並發更新 | Lock,或單個變數上的 atomic | Check-then-act、read-modify-write |
| Coordination | Threads 需要交接或等待 | Bounded blocking queue | 異步工作、workers、backpressure |
| Scarcity | 資源有固定上限 | Semaphore,或把 queue 當 pool 用 | 並發操作上限、connection pool |
多數設計從 correctness 開始。一旦存在 shared state,或吞吐量上來,coordination 與 scarcity 就會出現。真實系統常常三者並存。一次只給其中一個命名。
1. 兩條 thread,一個座位
一個進程內的訂座服務。瀏覽座位,訂下一個。第一稿看起來沒有問題。
public boolean bookSeat(String seatId, String userId) {
Seat seat = seats.get(seatId);
if (seat.isAvailable()) {
seat.book(userId);
return true;
}
return false;
}Alice 和 Bob 都想要 7A。Alice 檢查:空著。Bob 檢查:仍空著 —— Alice 還沒寫完。兩人都訂了。Bob 覆蓋了 Alice。Alice 以為自己有座位。她沒有。
Alice: is 7A free? yes
Bob: is 7A free? yes ← Alice has not booked yet
Alice: book 7A
Bob: book 7A ← overwrites AliceBug 就住在這段空隙裡。在 single-threaded 的 main thread 上,同一段函式是安全的。在多 thread 的 process 裡則不是。
Failure: 把 check-then-book 拆成兩步就上線,只因為它通過了單 thread 的測試。
2. Correctness
兩條 thread 同時碰同一份 shared state,狀態就會被寫壞。兩種形態幾乎覆蓋所有情況。
Check-then-act
先檢查一個條件,再據此行動。空隙裡,另一條 thread 可以改掉那個條件。
- 座位。 檢查是否空閒,再預訂。
- 停車場。 檢查車位是否空,再把車分配進去。
- Rate limiter。 檢查用戶是否仍在限額內,再放行請求。
- 庫存。 檢查存量,再扣掉一件。
檢查與行動必須是 一次 atomic operation。預設工具是 lock:臨界區內只有一條 thread,其餘等待。Java 的 synchronized 在整個 block 上持有 lock。Go 與 Rust 稱之為 mutex。Python 用 threading.Lock。
public synchronized boolean bookSeat(String seatId, String userId) {
Seat seat = seats.get(seatId);
if (seat.isAvailable()) {
seat.book(userId);
return true;
}
return false;
}Alice 拿到 lock,檢查,預訂,再釋放。Bob 等待。輪到他檢查時,7A 已被佔用。他得到的是 no。
Failure: 只 lock 了檢查、沒有 lock 預訂;或者每條 request 各拿一把 lock,座位 map 本身卻仍無保護地共享。
Read-modify-write
讀出一個值,據此計算,再寫回去。count++ 看起來像一次操作。它其實是三步:讀、加一、寫。
兩條 thread 都讀到 5。都算出 6。都寫下 6。一次 increment 消失了。點擊計數、餘額、庫存數量、指標匯總 —— 凡是更新依賴於當前值的,都會落到這個形態。
對 單個變數,atomic 能在一個 CPU 步驟裡做完 read-modify-write。現代晶片有 compare-and-swap。Java 把它暴露為 AtomicInteger。
AtomicInteger count = new AtomicInteger(10);
int next = count.incrementAndGet(); // 11 — hardware CAS, not count++Python 沒有內建 atomics。同樣的 increment 需要一把 lock。
lock = threading.Lock()
counter = 0
with lock:
counter += 1一個計數器、一個 flag、一項統計,用 atomic。一旦兩個變數必須保持一致 —— 例如賬戶之間的轉賬 —— atomics 幫不上忙。用一把覆蓋兩次 write 的 lock。
Failure: 在共享的 int 上做 count++,或用兩個 AtomicInteger 去做必須 all-or-nothing 的轉賬。
3. Coordination
工作必須從一條 thread 交到另一條。Signup 應當立刻返回。歡迎郵件要花 500ms。API threads 把任務放進 queue。Worker threads 再取下來。
Queue 一旦存在,兩個問題就會出現。
Worker 如何知道有工作來了 —— 還不至於燒掉一顆核心? 在空 queue 上 while (true) 輪詢,會永遠佔著一顆 CPU。Sleep-then-poll 少浪費一些週期,但引入延遲:任務若落在 100ms sleep 剛開始時,就要等這次 sleep 結束。
真正想要的是:queue 為空時休眠,producer 一 put 立刻醒來。這就是 blocking queue。Worker 調用 take。空則這條 thread 休眠。一次 put 會喚醒等待者。
BlockingQueue<Email> emails = new LinkedBlockingQueue<>(1000);
public void signup(User user) {
emails.put(new Email(user)); // blocks if the queue is full
}
// worker
while (true) {
Email task = emails.take(); // sleeps if empty
send(task);
}Python 的 queue.Queue 預設就是 blocking。put 與 get 是同一套想法。
若工作到達的速度超過 workers 的消耗呢? 無界 queue 會一直漲,直到 process 耗盡記憶體。給 queue 設上限。滿了之後,put 阻塞。這就是 backpressure:consumers 跟不上時,producers 自然放慢。始終設上限。
凡是工作在 threads 之間流動,就會出現這種形態:scheduler、後台 job processor、process 內部的消息交接。
本文不展開的後續問題:buffer 該有多大、shutdown 時 workers 如何排空、以及阻塞 request path 不可接受時該怎麼辦。
Failure: 在空 queue 上自旋,或在註冊高峰裡用一條無界 queue 吃掉 heap。
4. Scarcity
資源有固定上限。外部 API 只允許 10 個 in-flight 調用。五十條 thread 都想進去。需要一種說法:同一時刻只許十個。
Semaphore 是一桶 permits。做事前先取一張。做完再放回。桶空了:等到有人把 permit 放回來。
Semaphore downloads = new Semaphore(5);
public void download() throws InterruptedException {
downloads.acquire();
try {
doDownload();
} finally {
downloads.release();
}
}若 doDownload 拋錯而你從未 release,那張 permit 就丟了。五次異常會把桶清空。之後每一次 download 都會永遠等下去。在 finally 裡釋放。Python 是同一形態:acquire、做事、在 finally 裡 release。
有時你數的不是次數。你在復用 帶狀態的對象 —— 一條 database connection 握著 socket、記憶體、事務狀態。不會為每次 query 新開一條。你先建好十條 connection,再把它們遞出去。
這個 pool 就是裝滿這些對象的 blocking queue。take 一條,用完,再 put 回去。與 coordination 是同一種 primitive,只是當作一隻袋子來用。
BlockingQueue<Connection> pool = new LinkedBlockingQueue<>(10);
void init() throws InterruptedException {
for (int i = 0; i < 10; i++) {
pool.put(openConnection());
}
}
void query(String sql) throws InterruptedException {
Connection conn = pool.take();
try {
conn.execute(sql);
} finally {
pool.put(conn);
}
}兩種 scarcity 都是 acquire、使用、release。工作失敗時,release 仍然必須發生。
Failure: 只在成功路徑上歸還 permit 或 connection,第五次異常就會把整個 process 卡住。
5. 語言對照表
同一組 primitives,不同的名字。伸手去拿已經存在的那一個。
| Concept | Java | Python | Go | C++ | C# |
|---|---|---|---|---|---|
| Lock / mutex | synchronized / ReentrantLock | threading.Lock | sync.Mutex | std::mutex | lock / Monitor |
| Read-write lock | ReentrantReadWriteLock | N/A(第三方) | sync.RWMutex | std::shared_mutex | ReaderWriterLockSlim |
| Condition variable | Object.wait / notify | threading.Condition | sync.Cond | std::condition_variable | Monitor.Wait / Pulse |
| Semaphore | Semaphore | threading.Semaphore | x/sync/semaphore | std::counting_semaphore | SemaphoreSlim |
| Blocking queue | LinkedBlockingQueue | queue.Queue | buffered channel | 自行組合 | BlockingCollection |
| Atomic integer | AtomicInteger | N/A(用 lock) | sync/atomic | std::atomic | Interlocked |
| Concurrent map | ConcurrentHashMap | N/A(GIL) | sync.Map | TBB hash map | ConcurrentDictionary |
- Python 的 GIL 表示 CPU-bound 的 threads 不會並行。I/O-bound 的仍然會。沒有原生 atomic integer。
- Go:channel 是慣用的 queue。先用它,再考慮
sync.Cond。 - C++ 常常要你自己把 mutex 與 condition variable 組合起來。
6. 清單
三個問題。多數設計會落到其中之一。
- 是否存在多於一條 thread 能碰到的 shared state? Correctness。Check-then-act 或 read-modify-write。Lock 住臨界區;若只是單個變數,用 atomic。
- 工作是否從一條 thread 流向另一條? Coordination。空則休眠,滿則阻塞。一條 bounded blocking queue。
- 是否存在固定上限? Scarcity。在計數時用 semaphore。在 pooling 對象時用裝滿對象的 blocking queue。在
finally裡釋放。
先給類型命名,再給 primitive 命名,然後寫下臨界區。在這三種撐不住設計之前,不要發明第四種機制。
Failure: 還沒說清自己在解這三類裡的哪一類,就跳進自製的 lock-free 結構,或給整個服務包上一把全局 lock。