Work Queues
Pattern ที่เจอบ่อยที่สุด
หัวข้อที่มีชื่อว่า “Pattern ที่เจอบ่อยที่สุด”work queue (หรือเรียก task queue / competing consumers) คือม้างานของ RabbitMQ producer หนึ่งตัว publish งาน — resize รูปนี้, ส่ง email นี้, สร้าง report นี้ — แล้ว pool ของ worker แชร์งานกัน แต่ละงานไปหา worker ตัวเดียว
นี่คือวิธี scale background job งานเยอะไป? เพิ่ม worker บน queue เดิม worker เหล่านั้นแข่งกันดึง message และ RabbitMQ ส่งแต่ละ message ให้ตัวใดตัวหนึ่ง
flowchart LR p["Producer (publishes tasks)"] --> q["Queue: tasks"] q --> w1["Worker 1"] q --> w2["Worker 2"] q --> w3["Worker 3"]
Round-robin และจุดอ่อนของตัวเอง
หัวข้อที่มีชื่อว่า “Round-robin และจุดอ่อนของตัวเอง”โดย default RabbitMQ ส่ง message แบบ round-robin: worker 1, worker 2, worker 3, worker 1 วนไป โดยไม่สนว่าแต่ละ worker ยุ่งแค่ไหน ซึ่งโอเคถ้าทุกงานใช้เวลาเท่ากัน แต่ถ้า worker 1 ได้งานช้าติด ๆ กัน ส่วน worker 2 ได้งานเร็ว round-robin ก็ยังป้อนงานเท่ากัน worker 1 เลยตกงานค้าง ส่วน worker 2 นั่งว่าง
ตัวแก้คือ prefetch (basic.qos): บอก RabbitMQ ว่า “อย่าเพิ่งให้ message ใหม่กับ worker จนกว่าจะ ack ตัวที่แล้ว” ตั้ง prefetch เป็น 1 worker จะถือ unacked task ได้ทีละตัว worker ที่ช้าก็จะดึงงานถัดไปช้าลงเอง worker ที่ยุ่งจึงได้งานน้อยลงตามธรรมชาติ วิธีนี้เปลี่ยน round-robin ให้เป็น fair dispatch
const channel = await conn.createChannel();await channel.assertQueue('tasks', { durable: true });
// Fair dispatch: at most one unacked message per worker.await channel.prefetch(1);
channel.consume('tasks', async (msg) => { if (!msg) return; await doWork(msg.content.toString()); // process the task channel.ack(msg); // only now ask for the next}, { noAck: false });channel = conn.channel()channel.queue_declare(queue="tasks", durable=True)
# Fair dispatch: at most one unacked message per worker.channel.basic_qos(prefetch_count=1)
def on_task(ch, method, properties, body): do_work(body.decode()) # process the task ch.basic_ack(delivery_tag=method.delivery_tag) # ask for the next
channel.basic_consume(queue="tasks", on_message_callback=on_task)channel.start_consuming()ch, _ := conn.Channel()ch.QueueDeclare("tasks", true, false, false, false, nil)
// Fair dispatch: at most one unacked message per worker.ch.Qos(1, 0, false)
msgs, _ := ch.Consume("tasks", "", false, false, false, false, nil)for d := range msgs { doWork(string(d.Body)) // process the task d.Ack(false) // only now ask for the next}ack ทำให้ worker crash ได้อย่างปลอดภัย
หัวข้อที่มีชื่อว่า “ack ทำให้ worker crash ได้อย่างปลอดภัย”สังเกต noAck: false / manual ack ข้างบน นี่คือครึ่งหลังของ work queue ที่เชื่อถือได้ เพราะ worker ack หลัง งานเสร็จเท่านั้น worker ที่ crash กลางคันจึงยังไม่ได้ ack — RabbitMQ เลย redeliver งานนั้นให้ worker ตัวอื่น ไม่มีอะไรหายเงียบ ๆ
จับคู่กับ durable queue และ persistent message (จากโมดูล reliability) แล้ว task queue ของคุณจะรอดทั้ง worker crash และ broker restart
การจูน prefetch
หัวข้อที่มีชื่อว่า “การจูน prefetch”prefetch เป็น 1 ปลอดภัยสุดและให้ dispatch ที่แฟร์สุด แต่เพิ่ม round-trip ต่อ message สำหรับงานเร็วและสม่ำเสมอ prefetch สูงขึ้น (เช่น 10–50) จะทำให้ worker มีงานป้อนตลอดและ throughput สูง trade-off คือ:
| Prefetch | ผล |
|---|---|
| 1 | dispatch แฟร์สุด, throughput ต่ำสุด, ปลอดภัยสุดสำหรับงานช้า/ไม่สม่ำเสมอ |
| ปานกลาง (10–50) | throughput ดีสำหรับงานเร็ว, แฟร์น้อยลงบ้าง |
| ไม่จำกัด (0) | throughput สูงสุด แต่ worker ตัวเดียวอาจกวาด backlog ไปหมดจน worker อื่นอด |