Prefetch & QoS
ปัญหา: worker ตัวเดียวกินหมด
หัวข้อที่มีชื่อว่า “ปัญหา: worker ตัวเดียวกินหมด”เอา consumer สองตัวมาต่อกับ queue หนึ่ง RabbitMQ จะแจก message แบบ round-robin — ตัวละหนึ่งสลับกันไป ฟังดูยุติธรรม แต่ยุติธรรมโดย จำนวน ไม่ใช่โดย ปริมาณงาน ถ้า broker ดัน message ให้เร็วที่สุดโดยไม่คิด consumer ที่คว้า message ช้า ๆ ไปสิบตัวก็จมงาน ขณะที่อีกตัวนั่งว่าง ที่แย่กว่านั้น ด้วย manual ack และไม่มี limit broker อาจยัด message ที่ยัง unacked เป็น ร้อย ๆ ตัวให้ consumer ตัวเดียว จน memory ของ consumer เองระเบิด
ทางแก้คือ prefetch ตั้งผ่าน basic.qos
prefetch ทำอะไร
หัวข้อที่มีชื่อว่า “prefetch ทำอะไร”prefetch จำกัดจำนวน message ที่ยัง unacknowledged ที่ consumer ถือได้พร้อมกัน ตั้งเป็น 1 แล้ว RabbitMQ จะไม่ให้ message ตัวที่สองจนกว่าจะ ack ตัวแรก ตั้งเป็น 50 แล้ว consumer มี in flight ได้ถึง 50 ตัว
นี่เปลี่ยน round-robin ให้เป็น fair dispatch: consumer ที่ยุ่งจะหยุดรับงานใหม่จนกว่าจะตามทัน message จึงไหลไปหาใครก็ตามที่ว่างจริง ๆ โดยธรรมชาติ
flowchart TB
subgraph rr["ไม่มี prefetch limit — round-robin โดยจำนวน"]
q1["queue"] --> w1["worker A (ยุ่ง,
จมกับงานช้า)"]
q1 --> w2["worker B (ว่าง,
ไม่มีอะไรถูก queue ให้)"]
end
subgraph fd["prefetch = 1 — fair dispatch"]
q2["queue"] --> w3["worker A
(ทีละหนึ่ง)"]
q2 --> w4["worker B
(รับตัวว่างถัดไป)"]
end การเลือกค่า
หัวข้อที่มีชื่อว่า “การเลือกค่า”มี trade-off:
- prefetch = 1 — ยุติธรรมที่สุด และปลอดภัยที่สุดสำหรับงานยาวหรือไม่สม่ำเสมอ ต้นทุนคือ latency นิดหน่อย: consumer รอ ack round-trip ก่อนได้ message ตัวถัดไป เหมาะกับงานช้า (image processing, ส่ง email)
- prefetch สูงขึ้น (เช่น 10–100) — throughput ดีกว่าสำหรับงานเร็วและสม่ำเสมอ เพราะ consumer มี buffer พร้อมเสมอ ไม่ต้องรอ delivery ตัวถัดไป ต้นทุนคือ distribution ไม่สม่ำเสมอเท่าและ memory ต่อ consumer มากขึ้น
default ที่ดีสำหรับ background job ทั่วไปคือ เลขน้อย ๆ (1 สำหรับงานช้า, 10–50 สำหรับงานเร็ว) ปรับจูนโดยดู queue depth และ consumer utilization อย่าปล่อยให้ไม่จำกัดใน production
ตั้ง prefetch
หัวข้อที่มีชื่อว่า “ตั้ง prefetch”// at most 1 unacked message per consumer at a timeawait channel.prefetch(1);channel.consume('orders', handler, { noAck: false });# at most 1 unacked message per consumer at a timechannel.basic_qos(prefetch_count=1)channel.basic_consume(queue="orders", on_message_callback=handler, auto_ack=False)// prefetchCount=1, prefetchSize=0, global=falsech.Qos(1, 0, false)msgs, _ := ch.Consume(q.Name, "", false, false, false, false, nil)