ข้ามไปยังเนื้อหา

Work Queues

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"]
queue เดียว, worker แข่งกันหลายตัว

โดย 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 });

สังเกต 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 เป็น 1 ปลอดภัยสุดและให้ dispatch ที่แฟร์สุด แต่เพิ่ม round-trip ต่อ message สำหรับงานเร็วและสม่ำเสมอ prefetch สูงขึ้น (เช่น 10–50) จะทำให้ worker มีงานป้อนตลอดและ throughput สูง trade-off คือ:

Prefetchผล
1dispatch แฟร์สุด, throughput ต่ำสุด, ปลอดภัยสุดสำหรับงานช้า/ไม่สม่ำเสมอ
ปานกลาง (10–50)throughput ดีสำหรับงานเร็ว, แฟร์น้อยลงบ้าง
ไม่จำกัด (0)throughput สูงสุด แต่ worker ตัวเดียวอาจกวาด backlog ไปหมดจน worker อื่นอด
work queue ให้อะไรกับเรา?
ทำไมต้องตั้ง prefetch (basic.qos) เป็นค่าต่ำ?
worker crash หลังรับงานแต่ก่อน ack เกิดอะไรขึ้น?
ข้อเสียของ prefetch แบบไม่จำกัดคืออะไร?