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

Designing a Resilient Consumer

บทเรียนเรื่อง reliability แต่ละบทให้ safety net มาทีละอย่าง consumer ระดับ production ต้องการทั้งหมดนั้นทำงานร่วมกัน บทสรุปนี้ประกอบทั้งหมดเข้าเป็นดีไซน์เดียวที่คุณปรับใช้ได้

consumer ที่เชื่อถือได้บน production มีคุณสมบัติเหล่านี้:

  1. durable queue + persistent messages — เพื่อไม่ให้ broker restart ทำให้งานหาย
  2. manual acknowledgement — ack หลัง งานสำเร็จเท่านั้น เพื่อให้ crash กลางทาง redeliver
  3. prefetch ที่เหมาะสมprefetch ที่มีขอบเขต (สัก 10–50) เพื่อไม่ให้ consumer ตัวเดียวถูกยัด message ที่ยังไม่ ack เป็นพัน
  4. idempotent processing — เพราะ delivery เป็นแบบ at-least-once message เดิมมาถึงได้สองครั้ง ทำให้ประมวลผลซ้ำสองครั้งแล้วไม่มีผลเสีย (dedupe ด้วย message id)
  5. retry with backoff แล้วค่อย park — เมื่อล้มเหลว ให้ nack โดยไม่ requeue ทันที; route ผ่าน retry queue แบบ TTL + dead-letter; หลัง N ครั้ง ส่งไป parking (dead) queue ให้คนดู
  6. connection recovery — reconnect อัตโนมัติเมื่อ connection ไป broker หลุด และ re-declare topology
  7. graceful shutdown — เมื่อได้ SIGTERM ให้หยุดรับ delivery ใหม่, ทำงาน in-flight ให้เสร็จ (และ ack ให้เรียบร้อย) แล้วค่อย close channel และ connection
  8. observability — log และปล่อย metric ของ message ที่ processed, retried และ parked เพื่อให้เห็น failure
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 คน)"]
การจัดการ message หนึ่งตัวของ resilient 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);
});

ถ้าคุณจำอะไรได้แค่สามอย่างจากทั้งคอร์สนี้ ให้เป็นสามอย่างนี้: ack หลังสำเร็จ ไม่ใช่ก่อน (นี่คือความต่างระหว่างงานหายกับไม่หาย); ทำ processing ให้ idempotent (เพราะ at-least-once การันตีว่ามี duplicate); และ อย่า requeue poison message ใน loop แน่น ๆ (ใช้ TTL + dead-letter retry และ parking queue) ที่เหลือคือรายละเอียดที่ต่อยอดจากสามข้อนี้

resilient consumer ควร acknowledge message เมื่อไร?
ทำไม processing ต้อง idempotent?
วิธีที่ปลอดภัยในการจัดการ message ที่ล้มเหลวซ้ำ ๆ คืออะไร?
graceful shutdown ควรทำอะไรเมื่อได้ SIGTERM?