Polling Publisher
คุณมี Transactional Outbox อยู่แล้ว และต้องการ relay มาช่วยส่งแถวออกไป Transaction Log Tailing มีประสิทธิภาพดีก็จริง แต่ลากทั้ง CDC stack ตามมาด้วย ทั้ง replication slot, connector และ log retention ซึ่งอาจเกินกว่าที่ทีมอยากดูแล โดยเฉพาะช่วงเริ่มต้น สิ่งที่คุณต้องการคือ relay ที่สร้างและรันได้ด้วยแค่แอปพลิเคชันกับ database ของตัวเอง
ถ้าไม่มีสิทธิ์เข้าถึง transaction log ทางเดียวที่จะรู้ว่ามีแถว outbox ใหม่คือถาม database ตรง ๆ แต่ loop แบบง่าย ๆ มีหลุมพรางอยู่ ถ้ามี relay สองตัวรันพร้อมกันก็อาจหยิบแถวเดียวกันซ้ำ อาจ publish แถวไปแล้วลืมทำเครื่องหมายว่าส่งแล้ว และถ้า poll ถี่เกินไปก็กระหน่ำ database จนอ่วม
ดังนั้นแรงต่าง ๆ คือ: คุณต้องการ relay ที่ง่ายที่สุดเท่าที่จะเป็นไปได้โดยไม่มี infrastructure เพิ่ม คุณยังต้อง publish ทุกแถวอย่างน้อยหนึ่งครั้งและหลีกเลี่ยงการ publish แถวสองครั้งจาก relay เดียว และคุณต้องการให้ต้นทุนและ latency ของการ polling อยู่ในระดับที่สมเหตุสมผล
วิธีแก้
หัวข้อที่มีชื่อว่า “วิธีแก้”คำตอบคือรัน polling publisher ที่เป็น background loop ที่ทำงานทุกช่วงเวลาคงที่ ดึงแถวที่ยังไม่ถูกส่งจาก outbox มาเป็นชุดโดยเอาเก่าสุดก่อน publish ทีละอันไปยัง broker แล้วทำเครื่องหมายว่าส่งแล้วเมื่อ broker ตอบรับ ให้ดึงด้วย row lock ที่ข้ามแถวซึ่งถูก lock อยู่ relay หลายตัวจะได้รันพร้อมกันอย่างปลอดภัยโดยไม่แย่งแถวเดียวกัน หลัง publish เสร็จก็ลบแถวทิ้งหรือตั้ง timestamp published_at ไว้ เพื่อไม่ให้ถูกหยิบขึ้นมาซ้ำ
การส่งยังเป็นแบบ at least once อยู่ ถ้า relay crash หลัง publish แต่ก่อนทำเครื่องหมาย poll รอบถัดไปจะ publish แถวนั้นซ้ำ ซึ่งรับได้ เพราะ Idempotent Consumer ที่ปลายทางจะดูดซับ duplicate ให้เอง
sequenceDiagram
participant R as Polling Relay
participant DB as Outbox Table
participant B as Message Broker
loop every interval
R->>DB: SELECT unsent rows FOR UPDATE SKIP LOCKED
DB-->>R: batch of rows
R->>B: publish each message
B-->>R: ack
R->>DB: mark rows as sent (or delete)
end ตัวอย่าง
หัวข้อที่มีชื่อว่า “ตัวอย่าง”loop จะจองแถวเป็นชุดด้วย FOR UPDATE SKIP LOCKED เพื่อไม่ให้ relay ที่รันพร้อมกันชนกัน จากนั้น publish ทีละแถว แล้วทำเครื่องหมายว่าส่งแล้วภายใน transaction เดียวกับที่จองไว้
async function pollOnce(pool: Pool, broker: Broker): Promise<void> { const client = await pool.connect(); try { await client.query('BEGIN'); const { rows } = await client.query( `SELECT id, event_type, payload FROM outbox WHERE published_at IS NULL ORDER BY created_at LIMIT 100 FOR UPDATE SKIP LOCKED`, );
for (const row of rows) { await broker.publish(row.event_type, row.payload, { key: row.id }); await client.query( 'UPDATE outbox SET published_at = now() WHERE id = $1', [row.id], ); } await client.query('COMMIT'); } catch (err) { await client.query('ROLLBACK'); throw err; } finally { client.release(); }}
setInterval(() => pollOnce(pool, broker).catch(console.error), 1000);import time
def poll_once(conn, broker) -> None: with conn: # one transaction: claim, publish, mark sent with conn.cursor() as cur: cur.execute( """SELECT id, event_type, payload FROM outbox WHERE published_at IS NULL ORDER BY created_at LIMIT 100 FOR UPDATE SKIP LOCKED""" ) rows = cur.fetchall() for row_id, event_type, payload in rows: broker.publish(event_type, payload, key=row_id) cur.execute( "UPDATE outbox SET published_at = now() WHERE id = %s", (row_id,), )
while True: poll_once(conn, broker) time.sleep(1)func pollOnce(ctx context.Context, db *sql.DB, broker Broker) error { tx, err := db.BeginTx(ctx, nil) if err != nil { return err } defer tx.Rollback()
rows, err := tx.QueryContext(ctx, `SELECT id, event_type, payload FROM outbox WHERE published_at IS NULL ORDER BY created_at LIMIT 100 FOR UPDATE SKIP LOCKED`) if err != nil { return err }
type msg struct{ id, eventType string; payload []byte } var batch []msg for rows.Next() { var m msg if err := rows.Scan(&m.id, &m.eventType, &m.payload); err != nil { rows.Close() return err } batch = append(batch, m) } rows.Close()
for _, m := range batch { if err := broker.Publish(m.eventType, m.payload, m.id); err != nil { return err } if _, err := tx.ExecContext(ctx, `UPDATE outbox SET published_at = now() WHERE id = $1`, m.id); err != nil { return err } } return tx.Commit()}use sqlx::{Pool, Postgres, Row};
async fn poll_once(pool: &Pool<Postgres>, broker: &Broker) -> Result<(), sqlx::Error> { let mut tx = pool.begin().await?; // claim, publish, mark sent in one tx
let rows = sqlx::query( "SELECT id, event_type, payload FROM outbox \ WHERE published_at IS NULL \ ORDER BY created_at LIMIT 100 \ FOR UPDATE SKIP LOCKED", ) .fetch_all(&mut *tx) .await?;
for row in rows { let id: uuid::Uuid = row.get("id"); let event_type: String = row.get("event_type"); let payload: serde_json::Value = row.get("payload");
broker.publish(&event_type, &payload, id).await?; sqlx::query("UPDATE outbox SET published_at = now() WHERE id = $1") .bind(id) .execute(&mut *tx) .await?; }
tx.commit().await}ผลลัพธ์ที่ตามมา
หัวข้อที่มีชื่อว่า “ผลลัพธ์ที่ตามมา”สิ่งที่คุณได้:
- ความเรียบง่าย ไม่มี CDC connector ไม่มี replication slot ไม่มี log retention ให้ปรับจูน — แค่ loop, query และ database connection ที่คุณมีอยู่แล้ว นี่มักเป็น relay แรกที่ถูกต้อง
- ความสามารถในการพกพา แพตเทิร์นนี้เป็น SQL ธรรมดา จึงใช้ได้กับฐานข้อมูลเชิงสัมพันธ์ใด ๆ และเข้าใจกับทดสอบได้ง่ายมาก
- concurrency ที่ปลอดภัย
FOR UPDATE SKIP LOCKEDทำให้ relay หลายตัวแบ่งงานกันได้โดยไม่ publish แถวเดียวกันสองครั้ง
สิ่งที่คุณต้องแลก:
- latency จากการ polling โดยเฉลี่ยข้อความต้องรอราวครึ่งหนึ่งของช่วง polling ก่อนถูก publish ลดช่วงให้สั้นลงก็ได้ latency ดีขึ้น แต่แลกมาด้วยภาระที่เพิ่ม
- ภาระที่ลงบน database ทุกรอบ poll รัน query ไม่ว่าจะมีงานหรือไม่ ตาราง outbox และ index จึงรับ traffic ตลอดเวลา ที่เป็น overhead ที่ log tailing เลี่ยงได้พอดี
- ยังคงเป็น at-least-once การ crash ระหว่าง publish และ mark-sent ทำให้ publish ซ้ำ ดังนั้น consumer ต้องเป็น idempotent
กฎที่ใช้ได้จริงคือ เริ่มจาก polling publisher ก่อน เพราะง่ายและไม่มี dependency เพิ่ม แล้วค่อยย้ายไป transaction log tailing เมื่อ latency หรือภาระจากการ polling กลายเป็นปัญหาจริง
เนื้อหาที่เกี่ยวข้อง
หัวข้อที่มีชื่อว่า “เนื้อหาที่เกี่ยวข้อง”- Transactional Outbox — ผลิตแถวที่ relay นี้ poll
- Transaction Log Tailing — relay ที่ latency ต่ำกว่าแต่ซับซ้อนกว่า
- Idempotent Consumer — จำเป็นเพราะ polling ส่งแบบ at least once
| ข้อดี | ข้อแลกเปลี่ยน |
|---|---|
| ไม่ต้องการ database feature พิเศษ — ทำงานกับ database ใด ๆ | polling เพิ่ม load บน database ตลอดเวลา |
| publish event reliable — เชื่อมโยงกับ database transaction | latency สูงกว่า transaction log tailing — ขึ้นกับ polling interval |
| simple กว่า transaction log tailing — ไม่ต้อง access DB internals | ต้องจัดการ outbox table ที่โตขึ้น — cleanup เป็นสิ่งจำเป็น |
| เหมาะกับ low-to-medium throughput | polling interval ที่ผิด — สั้นเกิน = load สูง, นานเกิน = latency สูง |
ข้อผิดพลาดที่พบบ่อย
หัวข้อที่มีชื่อว่า “ข้อผิดพลาดที่พบบ่อย”Outbox ที่ไม่ Cleanup — event ที่ publish แล้วยังค้างใน outbox table อาการ:
- outbox table ใหญ่ขึ้นเรื่อย ๆ จนกระทบ performance
- query ช้าเมื่อต้อง scan หา unpublished event
- เพิ่ม cleanup job ที่ลบ published event หลังจาก retention period
Polling ถี่เกิน — poll ทุก 100ms บน database ที่มี write heavy workload อาการ:
- database CPU สูงจาก polling
- เพิ่ม load ที่ไม่จำเป็นบน production database
- ปรับ interval ตาม latency tolerance — 1-5 วินาทีมักเพียงพอ
💡 ตัวอย่างจากของจริง
ระบบ Banking เก่า:
- ใช้ polling publisher เพราะ infrastructure ไม่รองรับ CDC (Change Data Capture)
- poll outbox table ทุก 30 วินาที — ยอมรับ latency ได้สำหรับ batch settlement
Debezium (alternative):
- polling publisher เป็น simpler alternative ก่อน adopt Debezium
- migration path: polling publisher → transaction log tailing เมื่อ throughput สูงขึ้น