สะพานเชื่อม 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 ต่อ
-
เปิด connection แยกต่างหาก ไม่ใช้ shared pool เพราะคำสั่ง
LISTENต้องผูกกับ connection เดิมตลอดอายุการใช้งาน -
LISTENทั้ง 5 channel แล้ว map ไปยัง queue ปลายทางตามตารางนี้pg_notify channel RabbitMQ queue attribute_changedattribute_changefriend_track_event_insertedfriend_track_event_triggerform_submittedform_submitted_triggercampaign_clickcampaign_click_triggerbooking_eventbooking_event_trigger -
เมื่อได้รับ notification ระบบจะ unmarshal payload แล้ว re-marshal ด้วย
encoding/jsonเพื่อให้ byte ตรงกับ canonical form ของ V8 จากนั้น publish เข้า queue ที่ map ไว้ แบบ verbatim ไม่มี envelope ครอบ -
หาก 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.goNew(cfg, pub, log)— สร้าง bridge โดย mapping ระหว่าง channel กับ queue มาจากconfig.Queues.*ไม่ได้ hardcode ไว้ในโค้ดRun(ctx),listenLoop(),waitLoop(),handleNotification()catchupFromOutbox(),sweepOutbox(),cleanupOutbox()backoff(),IsConnected(),DisconnectedAt()
cmd/worker/main.go—runPgListener()ซึ่งลงทะเบียน 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บน channelcampaign_clickพร้อมกับ INSERT ลงevent_outbox(ดู schema ได้ในเอกสารโปรเจกต์ database) - RabbitMQ — profile นี้ publish อย่างเดียว ไม่มีการ consume
- ปลายทางของ event — consumer ของทั้ง 5 queue อยู่ที่ profile
trigger-workerยกเว้นbooking_event_triggerที่อยู่ที่ profilecron-scheduler