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

Your First Queue

มาส่ง message หนึ่งตัวผ่าน RabbitMQ กัน เพื่อให้เรียบง่ายที่สุด เราจะใช้ default exchange — direct exchange ที่ไม่มีชื่อ ซึ่ง RabbitMQ ให้มาอัตโนมัติ โดย routing key ก็คือชื่อ queue นั่นเอง ทำให้เรา “publish ตรงเข้า queue” ได้โดยยังไม่ต้องประกาศ exchange (ทุกบทเรียนถัด ๆ ไปใช้ named exchange จริง)

ก่อนอื่น run broker บนเครื่องด้วย Docker:

Terminal window
docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:3-management

port 5672 คือ AMQP (สำหรับแอปคุณ); 15672 คือ management UI (เปิด http://localhost:15672 login ด้วย guest/guest)

producer connect, เปิด channel, ประกาศ queue (เพื่อให้แน่ใจว่า queue มีอยู่จริง) แล้ว publish message เข้าไป

import amqp from 'amqplib';
const conn = await amqp.connect('amqp://localhost');
const channel = await conn.createChannel();
const queue = 'hello';
await channel.assertQueue(queue, { durable: false });
channel.sendToQueue(queue, Buffer.from('Hello, RabbitMQ!'));
console.log('sent: Hello, RabbitMQ!');
await channel.close();
await conn.close();

สังเกต exchange ที่เป็น string ว่าง ("") — นั่นคือ default exchange และ routing key คือชื่อ queue

consumer connect, ประกาศ queue เดียวกัน (การประกาศเป็น idempotent — ทำทั้งสองฝั่งได้อย่างปลอดภัย) แล้ว subscribe callback ของตัวเองจะทำงานทุกครั้งที่มี message ถูก deliver

import amqp from 'amqplib';
const conn = await amqp.connect('amqp://localhost');
const channel = await conn.createChannel();
const queue = 'hello';
await channel.assertQueue(queue, { durable: false });
console.log('waiting for messages...');
channel.consume(queue, (msg) => {
if (msg) {
console.log('received:', msg.content.toString());
channel.ack(msg);
}
});
flowchart LR
  p["Producer
sendToQueue('hello')"] -->|default exchange| q["Queue: hello"]
  q -->|deliver| c["Consumer
print + ack"]
flow ของ hello-world

รายละเอียดสำคัญ แม้จะเล็กน้อย:

  1. ทั้งสองฝั่งประกาศ queue การประกาศเป็น idempotent และมั่นใจว่า queue มีอยู่ไม่ว่าใครจะ start ก่อน
  2. message ก็แค่ bytes RabbitMQ ไม่สนใจ format ของคุณ — string, JSON, protobuf อะไรก็ได้ การ serialize เป็นหน้าที่ของคุณ
  3. consumer ack หลัง process แล้ว เป็นการบอก RabbitMQ ว่า “เสร็จแล้ว” เพื่อให้ message ถูกลบ ถ้าข้ามอันนี้ไป คุณจะได้ message เดิมอีกครั้งตอน reconnect (ที่เป็นความปลอดภัยที่คุณอยากได้พอดี — เพิ่มเติมในบท acknowledgement)
"default exchange" ที่ใช้ในตัวอย่าง hello-world คืออะไร?
ทำไมทั้ง producer และ consumer ต้องประกาศ queue?
message body ต้องอยู่ใน format ไหน?
ack ของ consumer ทำอะไร?