Transactional Outbox
service ของคุณต้องบันทึกการเปลี่ยนแปลงทางธุรกิจ พร้อมกับบอกระบบส่วนที่เหลือให้รู้ด้วย อย่างที่บทนำของโมดูลชี้ให้เห็น การทำสองอย่างนี้เป็นการเขียนสองครั้งที่แยกจากกัน ครั้งหนึ่งลง database อีกครั้งลง broker จึงเปิดช่องให้ crash ทำให้เกิด ghost event หรือ event หาย สิ่งที่คุณต้องการคือวิธีที่ผูกชะตาข้อความไว้กับการเปลี่ยนแปลงทางธุรกิจ ถ้าการเปลี่ยนแปลง commit ข้อความก็ต้องถูกส่งในที่สุด ถ้า roll back ก็ต้องไม่มีข้อความเลย
คุณไม่สามารถนำ message broker เข้าร่วมใน database transaction ของคุณได้ และ two-phase commit ข้ามระบบทั้งสองก็เป็นสิ่งที่คุณตัดทิ้งไปแล้ว แต่คุณ สามารถ ควบคุมฐานข้อมูลของคุณเองได้อย่างสมบูรณ์ ภายใน local transaction เดียว คุณเขียนได้กี่แถวก็ได้ ลงกี่ตารางก็ได้ ตามต้องการ และทั้งหมดจะ commit หรือ roll back ด้วยกัน
แรงที่ขัดกันจึงเป็นแบบนี้ คุณต้องการ atomicity ระหว่าง “การเปลี่ยนแปลงทางธุรกิจเกิดขึ้นแล้ว” กับ “ข้อความที่บอกเรื่องนั้นมีอยู่จริง” แต่คุณมีระบบเดียวที่ให้ transaction ได้ คือ database และสุดท้ายข้อความก็ยังต้องออกจาก database ไปให้ถึง broker ให้ได้อยู่ดี
วิธีแก้
หัวข้อที่มีชื่อว่า “วิธีแก้”ให้เพิ่มตาราง outbox เข้าไปใน database ของ service เวลาที่มีการเปลี่ยนแปลงทางธุรกิจ ให้เขียนข้อความขาออกเป็นแถวในตาราง outbox ภายใน transaction เดียวกัน กับการเปลี่ยนแปลงนั้น เมื่อ insert ทั้งสองอยู่ใน local transaction เดียว ก็จะ commit พร้อมกันหรือไม่เกิดขึ้นเลย dual write จึงยุบเหลือการเขียน atomic ครั้งเดียว
แถวใน outbox ไม่ใช่การส่ง แต่เป็น เจตนา ที่บันทึกไว้อย่างคงทนว่าจะส่ง จากนั้นมี component แยกต่างหากชื่อ message relay คอยอ่านแถว outbox ที่ commit แล้วไป publish ต่อยัง broker พร้อมทำเครื่องหมายแต่ละแถวว่าส่งแล้วหรือลบทิ้งเมื่อ broker ตอบรับ ส่วนวิธีที่ relay อ่าน outbox เป็นเนื้อหาของอีกสองบทเรียนถัดไป คือ tail transaction log หรือ poll ตาราง
flowchart LR
subgraph Service
H[Handler]
end
subgraph DB[Service Database]
BT[(business table)]
OT[(outbox table)]
end
R[Message Relay]
B[(Message Broker)]
H -->|one local transaction| BT
H -->|same transaction| OT
R -->|read committed rows| OT
R -->|publish, then mark sent| B ตัวอย่าง
หัวข้อที่มีชื่อว่า “ตัวอย่าง”การเคลื่อนไหวหลักคือ transaction เดียวที่ insert แถวธุรกิจและแถว outbox ด้วยกัน payload ที่ถูก serialize และ metadata ของข้อความไปอยู่ใน outbox ไม่มีการ publish อะไรจาก code path นี้เลย
import { randomUUID } from 'node:crypto';
// pool is a node-postgres Pool; the whole function is one DB transaction.async function confirmOrder(orderId: string, customerId: string, amount: number) { const client = await pool.connect(); try { await client.query('BEGIN');
await client.query( 'UPDATE orders SET status = $1 WHERE id = $2', ['CONFIRMED', orderId], );
const payload = JSON.stringify({ orderId, customerId, amount }); await client.query( `INSERT INTO outbox (id, aggregate_type, aggregate_id, event_type, payload) VALUES ($1, $2, $3, $4, $5)`, [randomUUID(), 'Order', orderId, 'OrderConfirmed', payload], );
await client.query('COMMIT'); // business row + outbox row commit together } catch (err) { await client.query('ROLLBACK'); throw err; } finally { client.release(); }}import jsonimport uuid
# conn is a psycopg connection; the `with` block is one DB transaction.def confirm_order(conn, order_id: str, customer_id: str, amount: int) -> None: with conn: # commits on success, rolls back on exception with conn.cursor() as cur: cur.execute( "UPDATE orders SET status = %s WHERE id = %s", ("CONFIRMED", order_id), )
payload = json.dumps( {"orderId": order_id, "customerId": customer_id, "amount": amount} ) cur.execute( """INSERT INTO outbox (id, aggregate_type, aggregate_id, event_type, payload) VALUES (%s, %s, %s, %s, %s)""", (str(uuid.uuid4()), "Order", order_id, "OrderConfirmed", payload), ) # business row + outbox row commit together herefunc ConfirmOrder(ctx context.Context, db *sql.DB, orderID, customerID string, amount int64) error { tx, err := db.BeginTx(ctx, nil) if err != nil { return err } defer tx.Rollback() // no-op once committed
if _, err = tx.ExecContext(ctx, `UPDATE orders SET status = $1 WHERE id = $2`, "CONFIRMED", orderID); err != nil { return err }
payload, _ := json.Marshal(map[string]any{ "orderId": orderID, "customerId": customerID, "amount": amount, }) if _, err = tx.ExecContext(ctx, `INSERT INTO outbox (id, aggregate_type, aggregate_id, event_type, payload) VALUES ($1, $2, $3, $4, $5)`, uuid.NewString(), "Order", orderID, "OrderConfirmed", payload); err != nil { return err }
return tx.Commit() // business row + outbox row commit together}use sqlx::{Pool, Postgres};use uuid::Uuid;
async fn confirm_order( pool: &Pool<Postgres>, order_id: &str, customer_id: &str, amount: i64,) -> Result<(), sqlx::Error> { let mut tx = pool.begin().await?; // one local transaction
sqlx::query("UPDATE orders SET status = $1 WHERE id = $2") .bind("CONFIRMED") .bind(order_id) .execute(&mut *tx) .await?;
let payload = serde_json::json!({ "orderId": order_id, "customerId": customer_id, "amount": amount, }); sqlx::query( "INSERT INTO outbox (id, aggregate_type, aggregate_id, event_type, payload) \ VALUES ($1, $2, $3, $4, $5)", ) .bind(Uuid::new_v4()) .bind("Order") .bind(order_id) .bind("OrderConfirmed") .bind(payload) .execute(&mut *tx) .await?;
tx.commit().await // business row + outbox row commit together}ผลลัพธ์ที่ตามมา
หัวข้อที่มีชื่อว่า “ผลลัพธ์ที่ตามมา”สิ่งที่คุณได้:
- เจตนา publish ที่ atomic ข้อความและการเปลี่ยนแปลงทางธุรกิจใช้ transaction เดียวกัน ดังนั้นคุณจะไม่มีวันมีอย่างหนึ่งโดยไม่มีอีกอย่าง ปัญหา dual-write ถูกแก้ที่ต้นทาง
- ความเป็นอิสระจาก broker write path ไม่แตะ broker อีกต่อไป ดังนั้น broker ที่ล่มไม่สามารถทำให้การดำเนินงานทางธุรกิจล้มเหลวหรือช้าลงได้ — outbox เพียงสะสมแถวไว้จนกว่า relay จะตามทัน
- audit trail โดยธรรมชาติ outbox เป็นบันทึกที่คงทนถาวรของทุก event ที่เซอร์วิสตั้งใจจะปล่อยออกไป
สิ่งที่คุณต้องแลก:
- ได้ at-least-once ไม่ใช่ exactly-once relay อาจ crash หลัง publish แต่ก่อนทำเครื่องหมายแถวว่าส่งแล้ว รอบถัดไปจึง publish ซ้ำ ดังนั้น consumer ต้องเป็น idempotent ที่เป็นบทเรียนสุดท้ายของโมดูลนี้
- relay ที่ต้องสร้างและรัน ต้องมีใครสักคน drain outbox อีกสองบทเรียนถัดไปครอบคลุมสองวิธีในการทำ
- การดูแลรักษา outbox แถวที่ส่งแล้วต้องถูกตัดทิ้ง ไม่เช่นนั้นตารางจะโตขึ้นไม่สิ้นสุด
เนื้อหาที่เกี่ยวข้อง
หัวข้อที่มีชื่อว่า “เนื้อหาที่เกี่ยวข้อง”- Transaction Log Tailing — วิธีหนึ่งในการสร้าง relay โดยไม่ต้อง polling
- Polling Publisher — วิธีที่ง่ายกว่าในการสร้าง relay
- Idempotent Consumer — จำเป็นเพราะ outbox รับประกันเพียงการส่งแบบ at-least-once
| ข้อดี | ข้อแลกเปลี่ยน |
|---|---|
| guarantee ว่า event ถูกส่งถ้า transaction สำเร็จ — ไม่มี event หาย | ต้องเพิ่ม outbox table และ publisher component |
| ไม่มี distributed transaction ระหว่าง database และ message broker | เพิ่ม latency เล็กน้อย — event ไม่ได้ส่ง synchronous ทันที |
| ทำงานกับ message broker ใด ๆ | outbox table ต้อง cleanup อย่างสม่ำเสมอ |
| reliable foundation สำหรับ event-driven architecture | ต้องเลือก polling publisher หรือ transaction log tailing เป็น mechanism |
ข้อผิดพลาดที่พบบ่อย
หัวข้อที่มีชื่อว่า “ข้อผิดพลาดที่พบบ่อย”Outbox ที่ผสม concern — เก็บ event หลาย type ใน outbox เดียวโดยไม่มี routing อาการ:
- publisher ต้องรู้ว่า event แต่ละ type ส่งไป topic ไหน
- outbox กลายเป็น god table ที่ทุกอย่างรวมกัน
- แยก outbox ตาม aggregate root หรือ topic
ไม่มี Publisher ที่ Reliable — publisher crash แล้วไม่มีกลไก retry อาการ:
- event ค้างใน outbox นาน — consumer ไม่ได้รับ event
- ต้องมี retry mechanism และ dead letter handling สำหรับ publisher
💡 ตัวอย่างจากของจริง
Eventuate (Chris Richardson):
- framework สำหรับ Transactional Outbox pattern
- ผู้บัญญัติ pattern นี้เป็นส่วนหนึ่งของ microservices pattern language
Shopify:
- ใช้ Transactional Outbox สำหรับ order event
- เมื่อ order สำเร็จ → outbox record สร้างใน transaction เดียว → publisher ส่งต่อ
- ทำให้ payment, inventory, notification sync กันได้อย่างน่าเชื่อถือ