Designing a Resilient Consumer
ทุกอย่างพร้อมกัน
หัวข้อที่มีชื่อว่า “ทุกอย่างพร้อมกัน”บทเรียนเรื่อง reliability แต่ละบทให้ safety net มาทีละอย่าง consumer ระดับ production ต้องการทั้งหมดนั้นทำงานร่วมกัน บทสรุปนี้ประกอบทั้งหมดเข้าเป็นดีไซน์เดียวที่คุณปรับใช้ได้
เช็กลิสต์
หัวข้อที่มีชื่อว่า “เช็กลิสต์”consumer ที่เชื่อถือได้บน production มีคุณสมบัติเหล่านี้:
- durable queue + persistent messages — เพื่อไม่ให้ broker restart ทำให้งานหาย
- manual acknowledgement — ack หลัง งานสำเร็จเท่านั้น เพื่อให้ crash กลางทาง redeliver
- prefetch ที่เหมาะสม —
prefetchที่มีขอบเขต (สัก 10–50) เพื่อไม่ให้ consumer ตัวเดียวถูกยัด message ที่ยังไม่ ack เป็นพัน - idempotent processing — เพราะ delivery เป็นแบบ at-least-once message เดิมมาถึงได้สองครั้ง ทำให้ประมวลผลซ้ำสองครั้งแล้วไม่มีผลเสีย (dedupe ด้วย message id)
- retry with backoff แล้วค่อย park — เมื่อล้มเหลว ให้
nackโดยไม่ requeue ทันที; route ผ่าน retry queue แบบ TTL + dead-letter; หลัง N ครั้ง ส่งไป parking (dead) queue ให้คนดู - connection recovery — reconnect อัตโนมัติเมื่อ connection ไป broker หลุด และ re-declare topology
- graceful shutdown — เมื่อได้ SIGTERM ให้หยุดรับ delivery ใหม่, ทำงาน in-flight ให้เสร็จ (และ ack ให้เรียบร้อย) แล้วค่อย close channel และ connection
- observability — log และปล่อย metric ของ message ที่ processed, retried และ parked เพื่อให้เห็น failure
flow พร้อม net ครบทุกชั้น
หัวข้อที่มีชื่อว่า “flow พร้อม net ครบทุกชั้น”flowchart TB
deliver["message ถูก deliver
(ถูกจำกัดด้วย prefetch)"] --> dedupe{"เคยเห็น id นี้ไหม?"}
dedupe -->|เคย| ackdup["ack (ข้ามแบบ idempotent)"]
dedupe -->|ไม่เคย| work["process งาน"]
work -->|สำเร็จ| ack["ack"]
work -->|ล้มเหลว| retry{"attempts < N?"}
retry -->|ใช่| ttl["nack -> retry queue (TTL) -> กลับ main"]
retry -->|ไม่| park["route ไป parking queue
(alert คน)"] โครง consumer ที่ทนทาน
หัวข้อที่มีชื่อว่า “โครง consumer ที่ทนทาน”รายละเอียดต่างกันไปตาม client แต่รูปทรงเหมือนกันทุกที่: prefetch ที่มีขอบเขต, manual ack ตอนสำเร็จ, dead-letter เมื่อล้มเหลวซ้ำ ๆ และ shutdown hook
import amqp from 'amqplib';
const conn = await amqp.connect('amqp://localhost');const channel = await conn.createChannel();
await channel.assertQueue('tasks', { durable: true });channel.prefetch(20); // bounded in-flight work
channel.consume('tasks', async (msg) => { if (!msg) return; const id = msg.properties.messageId; try { if (await alreadyProcessed(id)) { channel.ack(msg); return; } // idempotent await doWork(msg.content); await markProcessed(id); channel.ack(msg); // ack only after success } catch (err) { const attempts = (msg.properties.headers?.['x-attempts'] ?? 0) + 1; if (attempts >= 5) channel.nack(msg, false, false); // -> DLX/parking else republishToRetryQueue(channel, msg, attempts); // TTL + DLX retry channel.ack(msg); // remove the original }});
process.on('SIGTERM', async () => { // graceful shutdown await channel.close(); await conn.close(); process.exit(0);});import pika, signal
conn = pika.BlockingConnection(pika.ConnectionParameters("localhost"))channel = conn.channel()channel.queue_declare(queue="tasks", durable=True)channel.basic_qos(prefetch_count=20) # bounded in-flight work
def on_message(ch, method, props, body): msg_id = props.message_id try: if already_processed(msg_id): # idempotent ch.basic_ack(method.delivery_tag); return do_work(body) mark_processed(msg_id) ch.basic_ack(method.delivery_tag) # ack only after success except Exception: attempts = (props.headers or {}).get("x-attempts", 0) + 1 if attempts >= 5: ch.basic_nack(method.delivery_tag, requeue=False) # -> DLX/parking else: republish_to_retry_queue(ch, body, props, attempts) # TTL + DLX ch.basic_ack(method.delivery_tag)
channel.basic_consume(queue="tasks", on_message_callback=on_message)
def shutdown(*_): # graceful shutdown channel.stop_consuming(); conn.close()signal.signal(signal.SIGTERM, shutdown)channel.start_consuming()conn, _ := amqp.Dial("amqp://guest:guest@localhost:5672/")defer conn.Close()ch, _ := conn.Channel()defer ch.Close()
ch.QueueDeclare("tasks", true, false, false, false, nil) // durablech.Qos(20, 0, false) // bounded prefetchmsgs, _ := ch.Consume("tasks", "", false, false, false, false, nil)
go func() { for d := range msgs { id, _ := d.MessageId, d.Headers if alreadyProcessed(id) { d.Ack(false); continue } // idempotent if err := doWork(d.Body); err != nil { attempts := attemptCount(d.Headers) + 1 if attempts >= 5 { d.Nack(false, false) // -> DLX/parking } else { republishToRetryQueue(ch, d, attempts) // TTL + DLX d.Ack(false) } continue } markProcessed(id) d.Ack(false) // ack after success }}()
<-ctx.Done() // graceful shutdown: stop, finish in-flight, then closeนิสัยที่สำคัญที่สุด
หัวข้อที่มีชื่อว่า “นิสัยที่สำคัญที่สุด”ถ้าคุณจำอะไรได้แค่สามอย่างจากทั้งคอร์สนี้ ให้เป็นสามอย่างนี้: ack หลังสำเร็จ ไม่ใช่ก่อน (นี่คือความต่างระหว่างงานหายกับไม่หาย); ทำ processing ให้ idempotent (เพราะ at-least-once การันตีว่ามี duplicate); และ อย่า requeue poison message ใน loop แน่น ๆ (ใช้ TTL + dead-letter retry และ parking queue) ที่เหลือคือรายละเอียดที่ต่อยอดจากสามข้อนี้