RabbitMQ Publisher และ Topology
ภาพรวม
webhook-go เป็น producer-only gateway กล่าวคือทุก endpoint ทำงานเบา ๆ แล้วโยนงานลง RabbitMQ ในโปรเจกต์นี้ไม่มี consumer ไม่มี cron และไม่มี background job ใด ๆ เลย consumer ทั้งหมด อยู่ที่ worker-go
รูปแบบการส่งข้อความที่ใช้คือ bare JSON กล่าวคือ body ของ message เป็น payload object ตรง ๆ
ไม่มี envelope wrapper แบบ {type, data, meta} ที่ template ของบริษัทใช้เป็นมาตรฐาน เหตุผลคือ
consumer ฝั่ง worker-go ซึ่ง port มาจาก NestJS parse object ดิบอยู่แล้ว จึงต้องรักษา
compatibility ไว้ ในโค้ดมี comment เตือนเรื่องนี้ไว้อย่างชัดเจน
Routing ใช้รูปแบบที่ง่ายที่สุด คือ exchange เป็นชนิด direct และ routing key เท่ากับชื่อคิวเสมอ
Business Flow
Topology ที่ประกาศตอน boot
cmd/api/main.go เรียก AMQP.DeclareTopology(...) หลังจาก publisher ออนไลน์แล้ว
หากล้มเหลวจะบันทึก log ระดับ warn แล้วทำงานต่อ เนื่องจากคิวอาจถูกสร้างไว้ล่วงหน้าโดยฝั่ง infra
- ประกาศ exchange หลัก
line_exchangeชนิดdirectแบบ durable - ประกาศ DLQ exchange
line_exchange_dlqชนิด direct แบบ durable - สำหรับทุกคิวในรายการ
assertQueueแบบ durable พร้อม argumentx-max-priority: 10,x-dead-letter-exchange: line_exchange_dlqและx-dead-letter-routing-keyเป็นชื่อคิวbindQueueคิวเข้ากับline_exchangeโดยใช้ routing key เป็นชื่อคิวassertQueueคิว.dlqที่คู่กันแบบ durablebindQueueคิว.dlqเข้ากับline_exchange_dlqโดยใช้ routing key เป็นชื่อคิวหลัก
คิวทั้งหมดที่ประกาศ และผู้ใช้งานจริง
| คิว | env | publisher ในโปรเจกต์นี้ |
|---|---|---|
line_webhook | RABBITMQ_QUEUE_LINE_WEBHOOK | LINE Webhook Gateway เส้นทาง default และ cache miss |
line_forward_webhook | RABBITMQ_QUEUE_LINE_FORWARD_WEBHOOK | ไม่มี — ประกาศไว้เฉย ๆ เพราะการ forward ทำผ่าน HTTP โดยตรง ดู Webhook Forward |
tracking_log | RABBITMQ_QUEUE_TRACKING_LOG | บันทึก Tracking Log |
mookept_audience | RABBITMQ_QUEUE_MOOKEPT_AUDIENCE_QUEUE | จัดการสมาชิก Audience เฉพาะ PUT |
mbox_callback | RABBITMQ_QUEUE_MBOX_CALLBACK | MBOX Chatwoot Callback |
mbox_handoff | RABBITMQ_QUEUE_MBOX_HANDOFF | MBOX Agent Handoff |
booking_notification | RABBITMQ_QUEUE_BOOKING_NOTIFICATION | Booking Cancel Postback |
ไฟล์ .env.example ยังมีค่า RABBITMQ_QUEUE_LINE_CHANGE_RICHMENU=line_change_richmenu อยู่
แต่ internal/config/api.go ไม่ได้อ่านค่านี้ ถือเป็นของค้างจาก service เดิม
การ publish ผ่าน Manager.PublishRaw
- แปลง payload เป็น bytes ด้วย
json.Marshal - เปิด channel ใหม่ต่อการส่งหนึ่งครั้ง เพื่อไม่ให้ channel ที่เสียหายทำให้ publisher ทั้งตัวค้าง
- เปิด publisher confirm และลงทะเบียน listener ก่อนที่จะ publish เพื่อปิดช่องว่างของ race condition
- ส่งข้อความด้วย
DeliveryMode: Persistent,ContentType: application/jsonและค่าPriorityซึ่งปกติเป็น 0 - รอ ack หากไม่ได้รับ ack จะคืน error และหาก context หมดเวลาจะคืน
ctx.Err() - ผู้เรียกแต่ละรายจัดการ error ต่างกัน
- เส้นทางแบบ fire-and-forget เช่น LINE gateway และ tracking จะบันทึก log อย่างเดียว ข้อความที่ส่งไม่สำเร็จจะสูญหาย
- เส้นทางแบบ synchronous เช่น mbox callback และ
PUT /api/audienceจะตอบ error กลับไปยัง client
การเชื่อมต่อ
AMQP_URLSรองรับการระบุหลาย URL สำหรับการหมุนเวียนแบบ HA จัดการโดยinternal/amqp/manager.go- การ publish ถูกป้องกันด้วย mutex เนื่องจาก AMQP channel ไม่ concurrency-safe
- publisher ถูกห่อเป็น
deps.Dependencyผ่านinternal/deps/amqppubจึงเข้าสู่ระบบ readiness, supervisor และ graceful shutdown โดยอัตโนมัติ ดู Health Check, Metrics และการเฝ้าระวัง Dependency - หาก
AMQP_URLSว่างdeps.AMQPจะเป็น nil และทุกการ publish จะคืนErrNotConfiguredแต่ service ยังคง boot และตอบ health check ได้ตามปกติ
ไฟล์และฟังก์ชันหลัก
| ไฟล์ | ของสำคัญ |
|---|---|
internal/amqp/publisher.go | Manager.PublishRaw(ctx, exchange, queue, payload, priority) — ช่องทางเดียวที่ทุกโมดูลใช้ |
internal/amqp/topology.go | type Topology และ Manager.DeclareTopology(t) |
internal/amqp/manager.go | connection manager แบบหลาย URL, Channel(), Ready(), การ reconnect |
internal/amqp/envelope.go | envelope {type,data,meta} ของ template ซึ่งไม่ได้ถูกใช้ในโปรเจกต์นี้ |
internal/deps/amqppub/amqppub.go | ห่อ publisher เป็น deps.Dependency (Name/Connect/Ping/Close) |
internal/config/api.go | struct RabbitMQ และเมธอด Queues() ที่กำหนดลำดับการประกาศ |
internal/config/api_load.go | อ่านค่าจาก env พร้อมค่า default ของทุกคิว |
cmd/api/main.go | เรียก DeclareTopology ตอน boot โดย guard ด้วยเงื่อนไข AMQP != nil && Ready() |
จุดเชื่อมต่อกับ Service อื่น
- worker-go — เป็น consumer ของทุกคิวข้างต้นและเป็นคู่สัญญาหลัก การเปลี่ยนรูปร่างของ payload ที่นี่ถือเป็น breaking change ต่อฝั่ง worker ทันที
- RabbitMQ cluster — คิวและ DLX อาจถูกสร้างล่วงหน้าโดยทีม infra หรือผ่าน CRD ซึ่งไม่มีปัญหา
เพราะ
DeclareTopologyเป็น idempotent จึงประกาศซ้ำได้ - DLQ — ทุกคิวมี mirror
.dlqคู่กัน แต่ฝั่ง publisher นี้ไม่เคยใช้งาน DLQ จะทำงานเมื่อ consumer ฝั่ง worker reject ข้อความหรือ retry จนครบจำนวน - ค่า environment ทั้งหมดขึ้นต้นด้วย
RABBITMQ_ดูรายการเต็มได้ที่.env.example