โครงสร้าง Queue RabbitMQ และกลไก Retry / Dead-letter
ภาพรวม
ชั้น internal/mq คือหัวใจของ worker ทั้งระบบ ทำหน้าที่ประกาศ topology ของ RabbitMQ, publish job, consume job และตัดสินใจว่า message ที่เกิด error ควรถูก retry หรือส่งเข้า dead-letter queue (DLQ)
โมเดล retry ที่ใช้ ไม่ใช่ กลไก x-death ของ broker แต่เป็น application-level republish ซึ่ง port มาจาก rabbitmq.decorator.ts ของ NestJS เดิม เพื่อให้พฤติกรรมเหมือนระบบเดิมทุกประการ
ระบบมีทั้งหมด 38 queue ในชุด topology และทุก queue จะมี DLQ คู่กันเสมอในรูปแบบ <queue>.dlq
Business Flow
การประกาศ topology
ขั้นตอนนี้ทำตอน boot ของทุก profile ที่ใช้ RabbitMQ
- ประกาศ exchange หลัก
line_exchange(typedirect, durable) และ exchange สำหรับ DLQ ชื่อline_exchange_dlq(direct, durable) - วนประกาศทุก queue ในรายการด้วย
QueueDeclareพร้อม argument ดังนี้x-dead-letter-exchange = line_exchange_dlqx-dead-letter-routing-keyเท่ากับชื่อ queue นั้นx-max-priority = 10(ปรับได้ด้วยRABBITMQ_MAX_PRIORITY)
- Bind queue เข้ากับ
line_exchangeโดยใช้ routing key เท่ากับชื่อ queue เสมอ - ประกาศ
<queue>.dlqแบบ durable แล้ว bind เข้ากับline_exchange_dlqด้วย routing key เดียวกับชื่อ queue ต้นทาง
การ publish
- Producer ทุกตัวเรียก
Publish(ctx, queue, body, priority)ซึ่งส่งเข้าline_exchangeโดยใช้ชื่อ queue เป็น routing key - ตัว message เป็น JSON ของ payload ล้วน ๆ ไม่มี envelope ครอบ เพื่อให้เข้ากันได้กับ producer ฝั่ง TypeScript
- ใช้ publisher confirm ร่วมกับ
DeliveryMode: Persistentและมี mutex ป้องกันการใช้งาน channel พร้อมกัน เนื่องจาก AMQP channel ไม่ concurrency-safe
การ consume พร้อม retry และ DLQ
- เปิด channel ใหม่ต่อ 1 queue ตั้งค่า
Qos(prefetch)(RABBITMQ_PREFETCHค่าเริ่มต้น 10) และใช้ manual ack - Delivery ถูกประมวลผลพร้อมกันได้สูงสุดเท่ากับค่า prefetch ผ่าน goroutine pool ทำให้งานที่ช้าใน message เดียวไม่ block ทั้ง queue ส่วนการ ack, nack และ republish จะถูก serialize ด้วย mutex ของ channel
- หาก handler คืนค่า
nilจะทำการAck - หาก handler คืนค่า error จะตัดสินด้วย
classify(err)- Transient — network error เช่น
ECONNRESET,ETIMEDOUT,ECONNREFUSED,EPIPE,EAI_AGAINหรือ HTTP status 429 และ 500 ขึ้นไป และx-retry-countยังน้อยกว่า 3 → หน่วงเวลาแบบกำลังสองคือ 1 วินาที, 4 วินาที, 9 วินาที แล้ว republish เข้า default exchange พร้อมเพิ่มx-retry-countอีก 1 จากนั้น ack message เดิม - Permanent หรือ retry ครบ 3 ครั้งแล้ว →
Nack(requeue=false)ทำให้ message ตกลง<queue>.dlq
- Transient — network error เช่น
- Panic ที่เกิดใน handler จะถูก recover และถือเป็น permanent error (dead-letter) โดย process ไม่ล้ม
- หาก republish ล้มเหลว ระบบจะ
Nack(requeue=false)แทนการ ack เพื่อไม่ให้ message หายไปเงียบ ๆ
ไฟล์และฟังก์ชันหลัก
internal/mq/topology.go—AllQueues()(38 queue),MonitoredDLQQueues(),AssertTopology()internal/mq/client.go—Client,New(),Connect(), reconnect supervisor และName()ที่คืนค่า"rabbitmq"internal/mq/publisher.go—Publish(),PublishJSON()internal/mq/consumer.go—RunConsumer(),drain(),handle(), ค่าคงที่maxRetryCount = 3และretryCountOf()internal/mq/registrar.go—Registrar.Register(queue, handler),Run(ctx)internal/mq/errclass.go—IsTransientError()ซึ่งเป็น port ของerror-classifier.tsinternal/mq/errors.go—mq.Permanent(err)และmq.Transient(err)สำหรับบังคับผลการจำแนก errorinternal/mq/json.go— marshal ให้ผลลัพธ์ตรงกับJSON.stringifyของ V8
รายชื่อ queue ทั้งหมด
routing key ของทุก queue เท่ากับชื่อ queue
email, create_audience, update_audience, delete_audience, line_forward_webhook,
line_broadcast_rich_message, line_multicast_rich_message, line_webhook, process_campaign,
line_change_richmenu, line_sync_follower_user, tracking_log, line_user_last_activity,
calculate_campaign_stat, calculate_rich_menu_stat, calculate_rich_menu_stat_item,
line_auto_response, line_import_csv_user, import_mapping_job, mookept_audience,
attribute_change, attribute_change_realtime, audience_membership_trigger,
friend_track_event_trigger, friend_track_ref_import, audience_refresh, action_execute,
scheduled_trigger_fire, scheduled_action_fire, form_submitted_trigger, campaign_click_trigger,
web_request_execute, message_received_trigger, knowledge_index, mbox_callback, mbox_handoff,
booking_notification, booking_event_trigger
import_mapping_job เป็น queue ที่เกิดขึ้นใหม่ในฝั่ง Go ไม่มีอยู่ในระบบ NestJS เดิม ส่วน message_received_trigger ถูก assert ใน topology แต่ ไม่ได้ ถูก monitor DLQ (ดูรายละเอียดในเอกสารระบบเฝ้าระวัง Dead-letter Queue)
จุดเชื่อมต่อกับ Service อื่น
- Producer ภายนอก
line-management-webhook-go— publish เข้าline_webhook,tracking_log,line_forward_webhookline-management-cms-api-go— publish งาน campaign, audience, import, knowledge index และ rich menuline-management-client-api-go— publish งาน tracking, form และ booking
- Producer ภายใน — profile
pg-listenerที่แปลงมาจาก pg_notify, scanner ต่าง ๆ ใน profilecron-schedulerและ consumer หลายตัวที่ publish งานต่อ เช่น trigger ที่ publish เข้าaction_execute - Config ที่เกี่ยวข้อง —
RABBITMQ_URL,RABBITMQ_PREFETCH,RABBITMQ_MAX_PRIORITY,RABBITMQ_EXCHANGE_NAME,RABBITMQ_DLQ_EXCHANGE_NAMEและกลุ่มRABBITMQ_QUEUE_*