Retry & Backoff
naive retry ที่เผา CPU ของคุณ
หัวข้อที่มีชื่อว่า “naive retry ที่เผา CPU ของคุณ”message ประมวลผลไม่สำเร็จ — downstream API ล่ม, ข้อมูล resolve ไม่ได้ชั่วคราว reaction ที่นึกออกทันทีคือ nack แบบ requeue ให้ message กลับเข้า queue แล้วลองใหม่
ทำแบบนั้นคุณจะได้ tight poison loop message กลับไปที่หน้า queue ทันที fail ทันที และ worker ของคุณก็หมุนเร็วที่สุดเท่าที่ทำได้ กระหน่ำ dependency ที่ล่มอยู่และกลบ message ที่ปกติ immediate requeue แทบไม่เคยเป็นสิ่งที่คุณต้องการสำหรับ failure จริง
คุณต้องการสองอย่างที่วิธี naive ขาด: delay ระหว่างการ retry และ limit ว่าจะลองกี่ครั้งก่อนยอมแพ้
delayed retry ด้วย TTL + dead-letter exchange
หัวข้อที่มีชื่อว่า “delayed retry ด้วย TTL + dead-letter exchange”RabbitMQ ไม่มีปุ่ม “retry ใน 30 วินาที” แบบ native แต่คุณสร้างเองได้ด้วยการเอาสอง feature ที่คุณรู้อยู่แล้วมาต่อกัน:
- retry queue ที่มี message TTL (เช่น 30s) และ ไม่มี consumer message นั่งรอที่นั่นแล้วหมดอายุ
- dead-letter exchange บน retry queue นั้นชี้ กลับ ไปที่ main queue ของคุณ
message ที่ fail ถูก publish ไป retry queue รอจน TTL หมด แล้วถูก dead-letter กลับมาที่ main queue เพื่อลองอีกครั้ง — เป็น retry แบบ delay โดย delay กำหนดด้วย TTL ไม่มีการหมุนถี่ ๆ
flowchart LR
main["main queue"] -->|process fails| check{"attempts < N?"}
check -->|yes| retry["retry queue
(TTL 30s, no consumer)"]
retry -->|TTL expires, DLX| main
check -->|no| park["parking queue
(dead — inspect by hand)"] นับจำนวนครั้ง แล้ว park
หัวข้อที่มีชื่อว่า “นับจำนวนครั้ง แล้ว park”delay อย่างเดียวไม่พอ — message ที่จะ ไม่มีวัน สำเร็จ (malformed, อ้างถึงข้อมูลที่ถูกลบไปแล้ว) จะ retry ไปตลอด ดังนั้นให้ track attempt count โดยทั่วไปเก็บใน message header (x-retry-count หรืออ่านจาก header x-death ที่ RabbitMQ เพิ่มให้ทุกครั้งที่ dead-letter) หลังจากครบ N ครั้ง หยุด retry แล้ว route message ไปที่ parking queue (dead-letter queue ที่ไม่มี consumer อัตโนมัติ) ที่คนหรือ alert เข้าไปตรวจได้
// Declare a retry queue that dead-letters back to the main exchange after TTL.await channel.assertQueue('tasks.retry', { durable: true, arguments: { 'x-message-ttl': 30000, // wait 30s 'x-dead-letter-exchange': '', // default exchange 'x-dead-letter-routing-key': 'tasks', // back to the main queue },});
channel.consume('tasks', (msg) => { const attempts = (msg.properties.headers?.['x-retry-count'] ?? 0) + 1; try { doWork(msg.content); channel.ack(msg); } catch (err) { if (attempts >= 5) { channel.sendToQueue('tasks.parking', msg.content, { persistent: true }); } else { channel.sendToQueue('tasks.retry', msg.content, { persistent: true, headers: { 'x-retry-count': attempts }, }); } channel.ack(msg); // remove the original; we've re-routed it }});channel.queue_declare(queue="tasks.retry", durable=True, arguments={ "x-message-ttl": 30000, # wait 30s "x-dead-letter-exchange": "", # default exchange "x-dead-letter-routing-key": "tasks", # back to the main queue})
def on_task(ch, method, props, body): attempts = (props.headers or {}).get("x-retry-count", 0) + 1 try: do_work(body) ch.basic_ack(method.delivery_tag) except Exception: if attempts >= 5: ch.basic_publish("", "tasks.parking", body, pika.BasicProperties(delivery_mode=2)) else: ch.basic_publish("", "tasks.retry", body, pika.BasicProperties(delivery_mode=2, headers={"x-retry-count": attempts})) ch.basic_ack(method.delivery_tag) # remove original; re-routedch.QueueDeclare("tasks.retry", true, false, false, false, amqp.Table{ "x-message-ttl": int32(30000), // wait 30s "x-dead-letter-exchange": "", // default exchange "x-dead-letter-routing-key": "tasks", // back to the main queue})
for d := range msgs { attempts := retryCount(d.Headers) + 1 if err := doWork(d.Body); err == nil { d.Ack(false) } else if attempts >= 5 { ch.PublishWithContext(ctx, "", "tasks.parking", false, false, amqp.Publishing{DeliveryMode: amqp.Persistent, Body: d.Body}) d.Ack(false) } else { ch.PublishWithContext(ctx, "", "tasks.retry", false, false, amqp.Publishing{DeliveryMode: amqp.Persistent, Body: d.Body, Headers: amqp.Table{"x-retry-count": attempts}}) d.Ack(false) }}สำหรับ backoff แบบเพิ่มขึ้น (exponential) ให้ใช้ retry queue หลายตัวที่ TTL โตขึ้นเรื่อย ๆ — 10s, 1m, 10m — แล้วเลื่อน message ขึ้นบันไดตามจำนวนครั้งที่ retry เพิ่มขึ้น