Skip to main content

โครงสร้าง 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

  1. ประกาศ exchange หลัก line_exchange (type direct, durable) และ exchange สำหรับ DLQ ชื่อ line_exchange_dlq (direct, durable)
  2. วนประกาศทุก queue ในรายการด้วย QueueDeclare พร้อม argument ดังนี้
    • x-dead-letter-exchange = line_exchange_dlq
    • x-dead-letter-routing-key เท่ากับชื่อ queue นั้น
    • x-max-priority = 10 (ปรับได้ด้วย RABBITMQ_MAX_PRIORITY)
  3. Bind queue เข้ากับ line_exchange โดยใช้ routing key เท่ากับชื่อ queue เสมอ
  4. ประกาศ <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

  1. เปิด channel ใหม่ต่อ 1 queue ตั้งค่า Qos(prefetch) (RABBITMQ_PREFETCH ค่าเริ่มต้น 10) และใช้ manual ack
  2. Delivery ถูกประมวลผลพร้อมกันได้สูงสุดเท่ากับค่า prefetch ผ่าน goroutine pool ทำให้งานที่ช้าใน message เดียวไม่ block ทั้ง queue ส่วนการ ack, nack และ republish จะถูก serialize ด้วย mutex ของ channel
  3. หาก handler คืนค่า nil จะทำการ Ack
  4. หาก 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
  5. Panic ที่เกิดใน handler จะถูก recover และถือเป็น permanent error (dead-letter) โดย process ไม่ล้ม
  6. หาก republish ล้มเหลว ระบบจะ Nack(requeue=false) แทนการ ack เพื่อไม่ให้ message หายไปเงียบ ๆ

ไฟล์และฟังก์ชันหลัก

  • internal/mq/topology.goAllQueues() (38 queue), MonitoredDLQQueues(), AssertTopology()
  • internal/mq/client.goClient, New(), Connect(), reconnect supervisor และ Name() ที่คืนค่า "rabbitmq"
  • internal/mq/publisher.goPublish(), PublishJSON()
  • internal/mq/consumer.goRunConsumer(), drain(), handle(), ค่าคงที่ maxRetryCount = 3 และ retryCountOf()
  • internal/mq/registrar.goRegistrar.Register(queue, handler), Run(ctx)
  • internal/mq/errclass.goIsTransientError() ซึ่งเป็น port ของ error-classifier.ts
  • internal/mq/errors.gomq.Permanent(err) และ mq.Transient(err) สำหรับบังคับผลการจำแนก error
  • internal/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

note

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_webhook
    • line-management-cms-api-go — publish งาน campaign, audience, import, knowledge index และ rich menu
    • line-management-client-api-go — publish งาน tracking, form และ booking
  • Producer ภายใน — profile pg-listener ที่แปลงมาจาก pg_notify, scanner ต่าง ๆ ใน profile cron-scheduler และ consumer หลายตัวที่ publish งานต่อ เช่น trigger ที่ publish เข้า action_execute
  • Config ที่เกี่ยวข้องRABBITMQ_URL, RABBITMQ_PREFETCH, RABBITMQ_MAX_PRIORITY, RABBITMQ_EXCHANGE_NAME, RABBITMQ_DLQ_EXCHANGE_NAME และกลุ่ม RABBITMQ_QUEUE_*