Skip to main content

สะพานเชื่อม PostgreSQL NOTIFY ไปยัง RabbitMQ และ Event Outbox

ภาพรวม

Event จำนวนมากในระบบเกิดขึ้นที่ระดับฐานข้อมูล เช่น เมื่อมีผู้ใช้กดลิงก์ของ campaign ระบบจะ INSERT ลงตาราง tracking_log แล้ว Postgres trigger จะยิง pg_notify ออกมา

Profile pg-listener คือ process ที่เปิด connection พิเศษไว้ LISTEN ทั้งหมด 5 channel แล้วแปลงทุก notification เป็น message ส่งเข้า RabbitMQ ให้ trigger-worker ประมวลผลต่อ

เนื่องจาก pg_notify มีลักษณะ fire-and-forget หาก connection หลุดระหว่างทาง event จะสูญหาย ระบบจึงมีตาราง event_outbox ทำหน้าที่เป็นตัวสำรองแบบ durable โดย pg-listener จะทำ catch-up หลัง reconnect, sweep รายการที่ค้าง และ cleanup ข้อมูลเก่าเป็นระยะ

Business Flow

LISTEN แล้ว publish ต่อ

  1. เปิด connection แยกต่างหาก ไม่ใช้ shared pool เพราะคำสั่ง LISTEN ต้องผูกกับ connection เดิมตลอดอายุการใช้งาน

  2. LISTEN ทั้ง 5 channel แล้ว map ไปยัง queue ปลายทางตามตารางนี้

    pg_notify channelRabbitMQ queue
    attribute_changedattribute_change
    friend_track_event_insertedfriend_track_event_trigger
    form_submittedform_submitted_trigger
    campaign_clickcampaign_click_trigger
    booking_eventbooking_event_trigger
  3. เมื่อได้รับ notification ระบบจะ unmarshal payload แล้ว re-marshal ด้วย encoding/json เพื่อให้ byte ตรงกับ canonical form ของ V8 จากนั้น publish เข้า queue ที่ map ไว้ แบบ verbatim ไม่มี envelope ครอบ

  4. หาก connection หลุด ระบบจะ reconnect ด้วย exponential backoff min(2^n วินาที, 30s) สูงสุด 100 ครั้ง แล้วเรียก os.Exit(1) เพื่อให้ orchestrator restart pod ให้

การจัดการ Event Outbox

  • Catch-up — หลัง reconnect สำเร็จ ระบบจะดึงแถวใน event_outbox ที่เกิดขึ้นหลังเวลาที่ขาดการเชื่อมต่อ โดยใช้ FOR UPDATE SKIP LOCKED แล้ว publish ซ้ำ เพื่อกัน event หายในช่วง downtime
  • Sweep — ทุก 5 นาที ระบบกวาดแถวที่ค้างเกิน 30 วินาทีแล้วยังไม่ถูกส่งออกไป
  • Cleanup — ทุก 1 ชั่วโมง ระบบลบแถวที่เก่ากว่า 24 ชั่วโมงทิ้ง

Health check

  • ลงทะเบียน readiness probe ชื่อ pg-listener หาก LISTEN connection หลุด endpoint /readyz จะ fail ทันที
  • เปิด GET /health บน health port (PG_LISTENER_HEALTH_PORT ค่าเริ่มต้น 6010) ซึ่งตอบกลับด้วยข้อมูล status, pg.connected, pg.disconnectedAt และ timestamp

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

  • internal/pglistener/service.go
    • New(cfg, pub, log) — สร้าง bridge โดย mapping ระหว่าง channel กับ queue มาจาก config.Queues.* ไม่ได้ hardcode ไว้ในโค้ด
    • Run(ctx), listenLoop(), waitLoop(), handleNotification()
    • catchupFromOutbox(), sweepOutbox(), cleanupOutbox()
    • backoff(), IsConnected(), DisconnectedAt()
  • cmd/worker/main.gorunPgListener() ซึ่งลงทะเบียน readiness probe และ handler ของ /health
  • Queue ที่ publish ออกไป ได้แก่ attribute_change, friend_track_event_trigger, form_submitted_trigger, campaign_click_trigger และ booking_event_trigger

จุดเชื่อมต่อกับ Service อื่น

  • PostgreSQL — ใช้ pgx.Conn เฉพาะตัวสำหรับ LISTEN และเข้าถึงตาราง event_outbox ด้วย SQL ที่คัดลอกมาแบบ verbatim
  • Postgres trigger ฝั่งฐานข้อมูล — เช่น tracking_log_campaign_click_notify ที่ยิง pg_notify บน channel campaign_click พร้อมกับ INSERT ลง event_outbox (ดู schema ได้ในเอกสารโปรเจกต์ database)
  • RabbitMQ — profile นี้ publish อย่างเดียว ไม่มีการ consume
  • ปลายทางของ event — consumer ของทั้ง 5 queue อยู่ที่ profile trigger-worker ยกเว้น booking_event_trigger ที่อยู่ที่ profile cron-scheduler