Concurrency — ThreadPool ให้บริการหลาย client ด้วย Arc<Mutex>/RwLock
บท1 ปิดท้ายด้วยหนี้ก้อนหนึ่งที่เราสัญญาไว้ว่าจะจ่ายคืน: server for stream in listener.incoming() รับ client ทีละราย — ตราบใดที่รายหนึ่งยังไม่วางสาย ที่เหลือต้องรอ client ที่ช้าเพียงคนเดียว block ทุกคน บทนี้จ่ายคืนหนี้นั้น: เปลี่ยน serve loop แบบเรียงคิวให้เป็น thread poolthread poolกลุ่ม worker thread คงที่ที่ดึงงานจาก channel มาทำ ที่กระจายการเชื่อมต่อไปให้ worker หลายตัว โดยทั้งหมด share store ก้อนเดียว ผ่าน Arc<Mutex<…>> หรือ Arc<RwLock<…>>
ในแผนที่แนวคิด storage engine นี่คือชั้น thread-per-request pool + concurrency control — จุดที่ตำราฐานข้อมูล (Petrov) และ final project ของหนังสือ The Rust Programming Language บรรจบกัน โครง ThreadPool ที่เราสร้างในบทนี้คือ code จากหนังสือ บท 21 ตรงๆ ส่วนวินัยการถือ lock ให้แคบคือสิ่งที่ตำราฐานข้อมูลเรียกว่า concurrency control #21 สอน Arc/Mutex/mpsc เบื้องต้นมาแล้ว — บทนี้จะไป ลึกกว่านั้น: RwLock, ขอบเขตของ lock, และ poisoning
code ลงมือของคอร์สนี้อยู่ใน repo kaen-kvstore (code ตัวอย่างกำลังจัดทำ) — ตลอด 8 บทเราสร้าง key-value store บนเครือข่ายที่กู้คืนจาก crash ได้ หนึ่งตัว ด้วย Rust std ล้วน (ไม่มี async/tokio, ไม่มี serde) บทนี้แทนที่ serve loop แบบ single-threaded ของบท1 ด้วย ThreadPool (หนังสือบท 21) ที่ share store ผ่าน Arc<Mutex>/Arc<RwLock> — persistence/durability/compaction จากบท2–4 ยังอยู่ครบ เราแค่เพิ่มความสามารถรับหลาย client พร้อมกัน
ทำไม ThreadPool ไม่ใช่ spawn ต่อ connection ดื้อๆ
หัวข้อที่มีชื่อว่า “ทำไม ThreadPool ไม่ใช่ spawn ต่อ connection ดื้อๆ”ทางที่ง่ายที่สุดในการรับหลาย client คือ thread::spawn หนึ่งตัวต่อหนึ่งการเชื่อมต่อ แต่ถ้าเปิดรับตรงๆ แบบนั้น client 10,000 รายก็คือ 10,000 thread — แต่ละตัวกิน stack เป็น MB บวกภาระ context switch เครื่องจะล้มด้วยภาระของตัวมันเอง thread pool คือกลุ่ม worker thread จำนวนคงที่ ที่ดึงงานจาก channel มาทำทีละชิ้น — เพดานของ thread ถูกกำหนดไว้ล่วงหน้า ทำให้ทรัพยากรคาดเดาได้ นี่คือโครงที่หนังสือ The Rust Programming Language บท 21 (§21.2) สร้างขึ้นเป็น final project และเราหยิบมาตรงๆ
หัวใจมีสามชิ้น: งานหนึ่งชิ้นคือ closure ที่ boxed ไว้ — type Job = Box<dyn FnOnce() + Send + 'static> (FnOnce เพราะรันครั้งเดียว, Send เพราะข้าม thread, 'static เพราะ worker อาจรันมันเมื่อไรก็ได้); Vec<Worker> เก็บ handle ของ thread; และ mpsc::Receiver ตัวเดียว ที่ worker ทุกตัวต้องแย่งกันดึงงาน จึงต้องห่อด้วย Arc<Mutex<Receiver>> (หลายเจ้าของ + กันเข้าถึงพร้อมกัน)
use std::sync::{Arc, Mutex, mpsc};use std::thread::{self, JoinHandle};
type Job = Box<dyn FnOnce() + Send + 'static>;
pub struct ThreadPool { workers: Vec<Worker>, sender: Option<mpsc::Sender<Job>>, // Option เพื่อ take() ตอน shutdown}
impl ThreadPool { pub fn new(size: usize) -> ThreadPool { assert!(size > 0); let (sender, receiver) = mpsc::channel(); let receiver = Arc::new(Mutex::new(receiver)); // Receiver ตัวเดียว share ทุก worker let mut workers = Vec::with_capacity(size); for _ in 0..size { workers.push(Worker::new(Arc::clone(&receiver))); } ThreadPool { workers, sender: Some(sender) } } pub fn execute<F>(&self, f: F) where F: FnOnce() + Send + 'static, { self.sender.as_ref().unwrap().send(Box::new(f)).unwrap(); }}กับดักตัวแรก: lock/recv ต้องเป็น expression เดียว
หัวข้อที่มีชื่อว่า “กับดักตัวแรก: lock/recv ต้องเป็น expression เดียว”worker แต่ละตัววน loop ดึงงานจาก channel ที่ share กัน — และตรงนี้มีกับดักที่หนังสือบท 21 เตือนไว้เป็นพาดหัว code ที่ ถูก ต้องให้ receiver.lock().unwrap().recv() เป็น expression เดียว:
struct Worker { thread: Option<JoinHandle<()>>,}impl Worker { fn new(receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker { let thread = thread::spawn(move || loop { // expression เดียว: MutexGuard ถูก drop ที่ `;` -> ปล่อย lock ก่อน job() รัน let message = receiver.lock().unwrap().recv(); match message { Ok(job) => job(), Err(_) => break, // sender ถูก drop -> ช่องปิด -> ออกจาก loop (graceful shutdown) } }); Worker { thread: Some(thread) } }}เหตุที่มันต้องเป็นบรรทัดเดียวคือเรื่องของ จังหวะที่ MutexGuard ถูก drop เมื่อ receiver.lock().unwrap().recv() จบเป็น statement ที่ ; ตัว temporary guard ที่ค้ำ .lock() ไว้จะถูก drop ทันที — lock ถูกปล่อยคืน ก่อน ที่ job() จะเริ่มทำงาน worker ตัวนี้จึงถือ lock แค่ชั่ววินาทีที่ดึงงานออกมา แล้ว worker ตัวอื่นก็ดึงงานถัดไปได้ขนานกันทันที
ถ้าเผลอเขียนเป็น while let Ok(job) = receiver.lock().unwrap().recv() { job(); } แทน — code compile ผ่านสวยงาม แต่กฎ scope ของ while let ทำให้ guard มีชีวิต ตลอดทั้ง body ของ loop รวมถึงระหว่าง job() รันด้วย ผลคือ worker ตัวที่ถือ lock อยู่จะกัน worker อื่นทุกตัวไว้จน job() เสร็จ — มี worker เดียวเท่านั้นที่ทำงานจริงตลอดเวลา thread pool กลายเป็น single-threaded อย่างเงียบๆ นี่คือกับดักคลาสสิกที่ compiler ไม่เตือน เพราะมันไม่ใช่ error ด้าน memory safety
share store: Arc<Mutex> กับ Arc<RwLock>
หัวข้อที่มีชื่อว่า “share store: Arc<Mutex> กับ Arc<RwLock>”ทีนี้ worker หลายตัวต้องแตะ store ก้อนเดียวกัน ทางเลือกมีสอง แบบแรกคือ Arc<Mutex<HashMap<…>>> — ตรงไปตรงมา ทุกการเข้าถึง (อ่านหรือเขียน) ผูกขาด lock ทีละคน แบบที่สองคือ RwLockRwLocklock ที่ให้ผู้อ่านหลายคนพร้อมกัน แต่ผู้เขียนต้องผูกขาด — lock ที่แยกผู้อ่านกับผู้เขียนออกจากกัน: .read() ให้ ผู้อ่านหลายคนถือ lock พร้อมกันได้ ตราบใดที่ไม่มีใครเขียน ส่วน .write() ต้องผูกขาดเดี่ยวๆ สำหรับ key-value store ที่ อ่านมากกว่าเขียน — ซึ่งเป็นภาระงานที่พบบ่อยที่สุด — RwLock ให้ผู้อ่านขนานกันได้จริง เป็นกำไรด้าน throughput ที่ Mutex ให้ไม่ได้ ส่วน Mutex ยังคงเป็นค่าปริยายที่ง่ายกว่าเมื่อสัดส่วนอ่าน/เขียนพอๆ กัน
โครง handler ต่อหนึ่งการเชื่อมต่อจึงเป็นแบบนี้ (store: Store ถูก Arc::clone เข้ามาต่อ connection ก่อนส่งเข้า pool.execute(move || …)):
use std::collections::HashMap;use std::io::{Read, Write};use std::net::TcpStream;use std::sync::{Arc, RwLock};
type Store = Arc<RwLock<HashMap<String, String>>>;
fn handle_connection(mut stream: TcpStream, store: Store) { let mut op = [0u8; 1]; while stream.read_exact(&mut op).is_ok() { let key = read_lp(&mut stream); // blocking I/O — ยังไม่ถือ lock let reply = match op[0] { 0 => store.read().unwrap().get(&key).cloned(), // อ่านแบบ share clone ค่าออกมา 1 => { let val = read_lp(&mut stream); // parse ให้เสร็จ "ก่อน" ถือ lock store.write().unwrap().insert(key, val); None } // เขียนผูกขาด scope แคบสุด _ => { store.write().unwrap().remove(&key); None } }; // guard ทุกตัวถูก drop ตรงนี้ ก่อนจะเขียนตอบกลับ socket ข้างล่าง match reply { Some(v) => { let _ = stream.write_all(v.as_bytes()); } // เขียน socket โดยไม่ถือ lock None => { let _ = stream.write_all(b"\0"); } } }}กับดักตัวจริง: ถือ guard คร่อม I/O ที่ block
หัวข้อที่มีชื่อว่า “กับดักตัวจริง: ถือ guard คร่อม I/O ที่ block”สังเกตวินัยใน handler ข้างบนให้ดี — มันคือหัวใจของบทนี้: อ่านและ parse ทุกอย่างที่ block ได้ นอก lock, ถือ lock แค่ตอนแตะ map ใน scope ที่แคบที่สุด, แล้ว clone ค่าออกมา เพื่อให้ guard ถูก drop ก่อนที่เราจะเขียนตอบกลับ socket
เหตุผลอยู่ที่ hook ที่ต้องจำ: code ที่ถือ guard คร่อม socket I/O compile ผ่านสะอาด borrow checker ไม่มีทางช่วยคุณ เพราะมันไม่ใช่ data race — มันคือ bug เชิงประสิทธิภาพ นี่คือ ❌ version ดิบที่ต้องหลีกเลี่ยง:
// ❌ version ดิบ: ถือ write guard คร่อมการเขียน socket ที่ blockfn handle_bad(mut stream: TcpStream, store: Store) { let mut op = [0u8; 1]; while stream.read_exact(&mut op).is_ok() { let mut map = store.write().unwrap(); // ถือ lock ผูกขาด... let value = map.entry("k".to_string()).or_default().clone(); // ...แล้วยังถืออยู่ตลอดเวลาที่ block รอ socket ซึ่งอาจนานแค่ไหนก็ได้: let _ = stream.write_all(value.as_bytes()); } // guard เพิ่งถูก drop ตรงนี้ — ระหว่างนั้น worker อื่นทุกตัวถูกกันไว้หมด}code นี้ compile ผ่าน exit 0 — เราตรวจแล้ว compiler เห็นว่าไม่มีการเข้าถึงหน่วยความจำที่ไม่ปลอดภัย จึงเงียบ แต่ผลลัพธ์ตอน runtime คือหายนะ: ถ้า client รายหนึ่งเน็ตช้าหรือหยุดอ่าน write_all จะ block ค้าง และเพราะ write guard ยังถูกถืออยู่ worker ทุกตัว ที่ต้องการแตะ store ก็ค้างตามไปด้วย — thread pool ทั้งกองกลายเป็นเรียงคิวหลัง client ช้าเพียงรายเดียว ซึ่งเป็นปัญหา เป๊ะ กับที่บท1 บอกว่าจะแก้ compiler พิสูจน์ความปลอดภัยเชิง memory/thread ให้ได้ แต่ ไม่พิสูจน์ความถูกต้องเชิง concurrency ระดับสถาปัตยกรรม ให้ — วินัยการถือ lock ให้แคบเป็นความรับผิดชอบของเรา ไม่ใช่ของ borrow checker
Poisoning: เมื่อ worker 1 panic
หัวข้อที่มีชื่อว่า “Poisoning: เมื่อ worker 1 panic”ถ้า worker ตัว1 panic ระหว่างถือ lock อยู่ Rust จะทำเครื่องหมาย lock นั้นว่า poisoned — เพราะข้อมูลข้างในอาจถูกแก้ค้างไว้ครึ่งๆ กลางๆ ครั้งต่อไปที่ใครเรียก .lock()/.write() จะได้ Err(PoisonError) กลับมา จุดที่ต้องรู้คือ ความไม่สมมาตรระหว่าง Mutex กับ RwLock: Mutex poison เมื่อมี panic ขณะถือ lock แบบใดก็ตาม ส่วน RwLock poison เฉพาะเมื่อ panic เกิดขณะถือ write lock เท่านั้น — panic ในฝั่งผู้อ่านไม่ทำให้ lock เป็นพิษ (ผู้อ่านไม่ได้แก้ข้อมูล จึงไม่ทิ้งสถานะครึ่งๆ กลางๆ)
สำหรับ server ที่ต้องอยู่รอด เรามักไม่อยากให้ worker ตัวเดียว panic แล้วล้มทั้งระบบตาม — invariant ของ HashMap ที่เป็น store ธรรมดามักรอด panic ได้อยู่แล้ว เราจึง กู้ lock ที่เป็นพิษด้วย unwrap_or_else(|e| e.into_inner()) ซึ่งดึงค่าข้างในกลับมาใช้ต่อ แทนที่จะ .unwrap() ที่จะ panic ซ้ำแล้วลาม:
use std::sync::{Arc, Mutex};use std::collections::HashMap;
// getter ที่ทน poison: ให้บริการต่อได้แม้ worker ตัวก่อนหน้าจะ panicfn get(store: &Arc<Mutex<HashMap<String, String>>>, key: &str) -> Option<String> { let map = store.lock().unwrap_or_else(|e| e.into_inner()); // กู้ lock ที่เป็นพิษ map.get(key).cloned() // clone ออก guard drop หลังจากนั้น}RwLock ให้ผู้อ่านขนานกันได้ก็จริง แต่มีมุมมืดสองข้อจากเอกสาร std ที่ต้องรู้: (1) ลำดับความสำคัญ reader/writer ขึ้นกับระบบปฏิบัติการ — มาตรฐานไม่รับประกันว่า writer ที่รออยู่จะได้คิวก่อน reader ที่มาทีหลัง ถ้า reader หลั่งไหลไม่หยุด writer อาจอดตาย (writer starvation) บนบาง platform (2) การขอ read lock ตัวที่สองบน thread เดิมขณะที่ writer กำลังรอคิวอยู่ อาจ deadlock ได้ (เอกสาร std ยกตัวอย่างไว้เอง) — อย่าถือ read guard ค้างแล้วขอ read เพิ่มบน thread เดียวกัน วินัยเดิมยังใช้ได้: ถือ lock ใน scope ที่แคบที่สุด แล้วปล่อย
Graceful shutdown ผ่าน Drop
หัวข้อที่มีชื่อว่า “Graceful shutdown ผ่าน Drop”ตอนปิด server เราอยากให้ worker ทุกตัวทำงานที่ค้างอยู่ให้จบก่อนจากไป ไม่ใช่โดนตัดกลางคัน กลไกคือ impl Drop for ThreadPool: ปิด channel ก่อน แล้วค่อย join — drop(self.sender.take()) ทำให้ Sender หายไป ช่องส่งงานปิด ทุก recv() ที่ worker รออยู่จึงคืน Err แล้วออกจาก loop; จากนั้นวน join() ทีละ worker ให้ thread จบสนิท (ทั้ง sender และ thread เป็น Option เพื่อ take() ค่าออกมาได้ระหว่าง drop):
impl Drop for ThreadPool { fn drop(&mut self) { drop(self.sender.take()); // ปิด channel -> ทุก recv() คืน Err for worker in &mut self.workers { if let Some(thread) = worker.thread.take() { thread.join().unwrap(); // รอ worker ทำงานค้างให้จบแล้วเก็บ thread } } }}รันจริงบน musl: สร้าง pool 4 worker share Arc<Mutex<HashMap<String, String>>> ป้อนงาน insert 12 ชิ้น แล้วปล่อยให้ pool drop (graceful shutdown) ก่อนอ่าน map:
OK: 12 keys across 4 workers, graceful shutdown cleanfirst="key-0" last="key-9"ครบ 12 key ไม่มีตกหล่น (ไม่มี lost update) และ pool.join คืนค่าโดยไม่ค้าง — worker ทั้งสี่แบ่งงานกันทำจริง แล้วปิดตัวสะอาดเมื่อ channel ถูกปิด
flowchart LR
L["TcpListener :incoming()"] -->|"execute(job)"| CH["mpsc channel<br/>(Arc Mutex Receiver)"]
CH --> W1["Worker 1"]
CH --> W2["Worker 2"]
CH --> W3["Worker N"]
W1 --> S["store ที่ share<br/>Arc RwLock HashMap"]
W2 --> S
W3 --> S
คำบรรยายภาพ: ThreadPool — listener ป้อนงานผ่าน channel เดียว (Receiver ห่อด้วย Arc<Mutex<…>>) ให้ worker จำนวนคงที่แย่งกันดึงไปทำ แต่ละ worker แตะ store ก้อนเดียวกันที่ share ผ่าน Arc<RwLock<HashMap>> โดยถือ lock ใน scope ที่แคบที่สุด
เส้นแบ่งที่เราพูดตรงๆ
หัวข้อที่มีชื่อว่า “เส้นแบ่งที่เราพูดตรงๆ”model thread-per-request + ThreadPool + Arc<Mutex/RwLock> นี้ ถูกต้องและง่าย — มันคือ final project ของหนังสือบท 21 ตรงๆ มันหยุดสเกลตอนชนกำแพง C10k (connection ที่ส่วนใหญ่ idle นับหมื่น) เพราะแต่ละ OS thread กิน stack เป็น MB บวกภาระ context switch หลักคิดคือ “thread เป็นค่าปริยายที่ถูกต้องจนกว่าจำนวน connection — ไม่ใช่ CPU — จะกลายเป็นคอขวด”
อีกจุดที่ต้องพูดตรง: mpsc::channel() เป็น channel แบบ ไม่จำกัดความจุ (unbounded) — ไม่มี backpressure ถ้างานเข้าเร็วกว่าที่ worker ทำได้ คิวจะบวมกินหน่วยความจำไปเรื่อยๆ channel เดียวใน std ที่มีขอบเขตคือ sync_channel(n) ซึ่งจะ block ฝั่งส่ง เมื่อคิวเต็ม (นั่นคือคือ backpressure) ส่วน backpressure แบบ async ที่แท้จริงเป็นเรื่องของ runtime แยก — เลื่อนไปคอร์ส #23 การ ตั้งชื่อเส้นแบ่ง นี้ให้ชัดคือบทเรียน ไม่ใช่ช่องโหว่
สรุปก่อนไปต่อ
หัวข้อที่มีชื่อว่า “สรุปก่อนไปต่อ”บทนี้จ่ายคืนหนี้จากบท1: เราสร้าง ThreadPool ของหนังสือบท 21 (Job = Box<dyn FnOnce() + Send + 'static>, Vec<Worker>, Receiver เดียวห่อ Arc<Mutex<…>>) โดยระวังกับดักว่า receiver.lock().unwrap().recv() ต้องเป็น expression เดียวเพื่อให้ guard drop ก่อน job() รัน; share store ด้วย Arc<Mutex> (ค่าปริยายที่ง่าย) หรือ Arc<RwLock> (ผู้อ่านหลายคนพร้อมกันเมื่ออ่านมากกว่าเขียน); ยึดวินัย lock scope — อ่าน/parse นอก lock แล้ว mutate ใต้ lock ใน scope แคบสุด เพราะ hook สำคัญคือการถือ guard คร่อม socket I/O compile ผ่านแต่ทำให้ทุก worker ต่อคิว compiler ไม่ช่วย; กู้ lock ที่เป็นพิษด้วย into_inner() (จำความไม่สมมาตร: RwLock poison เฉพาะ writer panic); และปิดสะอาดผ่าน Drop ที่ drop sender ก่อนแล้ว join ทุก worker ทุก snippet compile และรันได้จริงบน Rust 1.97.1 / edition 2024 / std ล้วน
บท6 เราพิสูจน์การกู้คืนจากการล่ม: ตอนนี้ store รับหลาย client และ durable แล้ว (บท3) — บทหน้าเราจะเปลี่ยนคำกล่าวอ้างเรื่อง durability ให้เป็น การทดสอบอัตโนมัติ ด้วย cargo test จำลอง crash จริงด้วย SIGKILL (ที่ไม่รัน destructor เลย) แล้ว reopen, replay, ยืนยันว่าได้สถานะเดิมเป๊ะ
บทนี้อิงต้นทางที่ลงวันที่กำกับ อ่านต่อได้โดยตรง:
- The Rust Programming Language — §21.2 “Turning Our Single-Threaded Server into a Multithreaded Server” (เข้าถึง 2026-07-24) — โครง
ThreadPool(Job = Box<dyn FnOnce()…>,Vec<Worker>,Receiverห่อArc<Mutex<…>>) และคำเตือนพาดหัวว่า lock/recv ต้องเป็น expression เดียว - The Rust Programming Language — §21.3 “Graceful Shutdown and Cleanup” (เข้าถึง 2026-07-24) —
Option<Sender>+Option<JoinHandle>,Dropที่ drop sender ก่อนแล้ว join ทุก worker - std
sync::RwLock(เข้าถึง 2026-07-24) — ผู้อ่านหลายคน/ผู้เขียนเดี่ยว, poison เฉพาะ writer panic, ลำดับความสำคัญขึ้นกับ OS, ตัวอย่าง reader-reacquire deadlock - std
sync::Mutex(เข้าถึง 2026-07-24) — poisoning และPoisonError::into_innerเพื่อกู้ lock ที่เป็นพิษ - std
sync::mpsc(เข้าถึง 2026-07-24) —channel()แบบ unbounded (ไม่มี backpressure) เทียบกับsync_channel(n)ที่ block ฝั่งส่ง
เช็กความเข้าใจ — บทที่ 5
ข้อ 1 / 3ทำไม receiver.lock().unwrap().recv() ในตัว worker ต้องเป็น expression เดียว (ไม่ใช่ while let Ok(job) = receiver.lock().unwrap().recv())?