Idempotency & Consistency
สิ่งที่จะสร้าง
หัวข้อที่มีชื่อว่า “สิ่งที่จะสร้าง”บทนี้ไม่มีโค้ดใหม่ The Saga Handler → เขียนทุกบรรทัดที่จะพูดถึงไว้ครบแล้ว — การ dedupe ด้วย insert into processed_events ... on conflict do nothing, guard terminal-state ด้วย select status ... for update และ pgx.Tx เดียวที่ครอบทั้งสองอย่างไว้กับการ update status และ outbox row แต่บทนั้นใช้ทั้งหมดโดยไม่ได้หยุดอธิบายว่า ทำไม แต่ละชิ้นถึงมีรูปร่างแบบนั้น
บทนี้ทำหน้าที่นั้นพอดี — ตั้งชื่อ guard สองตัวใน OrderRepo.ApplyPaymentResult อธิบายว่าแต่ละตัวมีไว้รอด failure แบบไหนภายใต้ at-least-once delivery ของ Kafka แล้วพูดตรง ๆ เรื่องการรับประกันสองอย่างที่ saga นี้จงใจ ไม่ ให้ คือ exactly-once delivery และ rollback อัตโนมัติเมื่อขั้นตอนใดล้มเหลว
รูปแบบเดียวกับที่ Acks, Retry & Dead Letters → ใช้กับ RabbitMQ: บทก่อนเขียนกลไกไว้แล้วและใช้งานไปเลยโดยไม่มีพิธีรีตอง ส่วนบทนี้คือการอ่านละเอียดที่เปลี่ยน “โค้ดทำงานได้” ให้เป็น “รู้แน่ว่าทำไมถึงทำงาน และหยุดทำงานตรงไหน”
เริ่มจากสิ่งที่บังคับทั้งหมดนี้: at-least-once delivery
Producer & Consumer → พูดไว้ว่า Consumer.Run จะ commit offset หลัง handler return เท่านั้น ดังนั้นถ้า saga consumer crash หลัง ApplyPaymentResult commit แต่ยังไม่ทัน commit offset ตอน restart Kafka จะ redeliver event payment.succeeded เดิมเป๊ะ
แปลว่า saga จะ เห็น event บางตัวมากกว่าหนึ่งครั้ง ไม่ใช่ edge case หายาก แต่เป็นการรับประกันของ delivery model ที่ระบบนี้ตั้งอยู่บน guard ทุกตัวข้างล่างมีอยู่เพราะข้อเท็จจริงเดียวนี้
Guard 1 — การ dedupe ด้วย processed_events สิ่งแรกที่ ApplyPaymentResult ทำใน transaction คือพยายามบันทึก id ของ event:
tag, err := tx.Exec(ctx, ` insert into processed_events (event_id) values ($1) on conflict do nothing`, eventID)if err != nil { return fmt.Errorf("repo: record processed event %s: %w", eventID, err)}if tag.RowsAffected() == 0 { // Already processed this exact event — commit the no-op and return // without touching the order at all. return tx.Commit(ctx)}table processed_events มี event_id text primary key โดย migration ของ The Saga Handler → สร้างไว้ให้แล้ว on conflict do nothing จึงแปลว่า “insert ถ้ายังไม่เคยเห็น id นี้” แบบ race-free
ถ้า row มีอยู่แล้ว แปลว่า event ตัวนี้เป๊ะ ๆ ผ่านการ process ไปแล้วในการ deliver ครั้งก่อน RowsAffected() จะเป็น 0 แล้ว handler ก็ commit และ return โดย ไม่แตะอะไรอย่างอื่นเลย ไม่มีการประเมิน status ของคำสั่งซื้อซ้ำ และไม่มีการเขียน outbox row ตัวที่สอง
idempotency key คือ eventID ซึ่งก็คือ e.ID ที่ saga ส่งตรงเข้ามาจาก The Saga Handler → และ guard นี้ยังกันการ redelivery ของ Payment เอง ได้ด้วย ไม่ใช่แค่ของ saga เพราะ Process & Publish → derive id ของ event payment.* แต่ละตัวแบบ deterministic เป็น e.ID + ":payment"
ดังนั้น order.created ที่ redeliver จะผลิต payment.succeeded ที่มี id เดิม เสมอ และ processed_events ก็รู้ทันทีว่าเคยเห็นแล้ว id ที่ deterministic ตั้งแต่ต้นทางคือสิ่งที่ทำให้ dedupe ที่ปลายทางเป็นไปได้
Guard 2 — การเช็ค terminal-state ด้วย for update การผ่าน dedupe มาได้แปลว่านี่คือ event id ที่ saga ไม่เคย process จริง ๆ แต่แค่นั้นยังไม่พอที่จะ apply แบบหลับหูหลับตา:
var status stringif err := tx.QueryRow(ctx, ` select status from orders where id = $1 for update`, orderID,).Scan(&status); err != nil { return fmt.Errorf("repo: lock order %s: %w", orderID, err)}if status != "pending" { // The order already left "pending" — a duplicate or out-of-order // payment event arrived after the saga already resolved it. Commit // the processed_events insert above and stop; the order's status // is not touched a second time. return tx.Commit(ctx)}ตรงนี้มีสองอย่างเกิดขึ้นพร้อมกัน
อย่างแรก for update จับ row-level lock บนคำสั่งซื้อนั้นไว้ตลอดที่เหลือของ transaction ถ้ามี payment-result event สองตัวของคำสั่งซื้อเดียวกันเข้ามาพร้อมกัน (saga consumer สอง instance) transaction ตัวที่สองจะ block อยู่ที่ SELECT นี้จนกว่าตัวแรกจะ commit จึงไม่มีทางอ่าน status เก่าแล้วตัดสินใจผิด
อย่างที่สอง การเช็ค if status != "pending" คือ terminal-state guard คำสั่งซื้อที่ถึง confirmed หรือ cancelled แล้วถือว่าจบ payment event ที่ตามมาทีหลัง — payment.failed ที่มาถึงหลัง payment.succeeded confirm ไปแล้ว หรือ event ที่ id ต่างกันจริง ๆ ของคำสั่งซื้อที่ event อื่น resolve ไปแล้ว — ต้องไม่ขยับสถานะเป็นครั้งที่สอง handler จะ commit เฉพาะการ insert processed_events เพื่อบันทึก id นั้นไว้ แล้วหยุด
ทำไมต้องมีสอง guard ในเมื่อฟังดูเหมือนตัวเดียวน่าจะพอ? เพราะสอง guard จับคนละเคส
dedupe จับ event เดิม ที่มาสองครั้ง — id เดียวกัน ส่วนการเช็ค terminal-state จับ event คนละตัว ที่พยายาม re-resolve คำสั่งซื้อที่ terminal ไปแล้ว — id ต่างกันแต่คำสั่งซื้อเดียวกัน ที่เป็นเคสที่ dedupe ยินดีปล่อยผ่าน
ระบบที่มีแค่ dedupe จะ ignore payment.succeeded ที่ redeliver ได้ถูกต้อง แต่ยังปล่อยให้ payment.failed ที่มาสายพลิกคำสั่งซื้อ confirmed กลับเป็น cancelled ได้ ส่วนระบบที่มีแค่การเช็ค terminal-state จะรอดเคสนั้น แต่จะรัน status-update พร้อมเขียน outbox ใหม่ทั้งชุดทุกครั้งที่ event เดิม redeliver จนกว่าจะมีตัวใดตัวหนึ่ง land ก่อน ซึ่งทั้งเปลืองและ racy เพราะไม่มี lock
พอวางสอง guard ไว้ด้วยกันใน transaction เดียว ApplyPaymentResult ก็ถูกต้องภายใต้ ordering, duplication หรือ concurrency แบบไหนก็ตามที่ delivery layer โยนมา
Transaction คือ guard ตัวที่สามแบบเงียบ ๆ ทุกอย่าง — dedupe insert, การอ่านแบบ locked, การ update status และ outbox row — อยู่ใน pgx.Tx เดียวพร้อม defer tx.Rollback(ctx) วินัยเดียวกับที่ The Order Repository → ใช้ทุกที่
ผลคือมีแค่สองทาง: คำสั่งซื้อขยับไป terminal status และ บันทึก processed_events row และ เขียน outbox row order.confirmed/order.cancelled สำเร็จพร้อมกันทั้งหมด หรือไม่เกิดอะไรขึ้นเลยสักอย่าง จึงไม่มีช่องว่างที่ event ติดสถานะ processed แต่ status ไม่ขยับ หรือ status ขยับแล้วแต่ไม่มี outbox row ไปบอกส่วนที่เหลือของระบบ
ทีนี้ความซื่อสัตย์ สองอย่างที่ saga นี้ ไม่ ให้ โดยตั้งใจ:
- ไม่ใช่ exactly-once delivery แต่เป็น at-least-once delivery พร้อม effect แบบ exactly-once guard พวกนี้ไม่ได้หยุด Kafka จากการ deliver event เดิมสิบครั้ง แต่ทำให้การ deliver ครั้งที่สิบกลายเป็น no-op ระบบ รับ event หลายครั้ง แต่ effect ต่อคำสั่งซื้อ เกิดขึ้นครั้งเดียว ความต่างตรงนี้สำคัญ เพราะ “exactly-once delivery” ในงาน distributed messaging ทั่วไปแทบทำไม่ได้เลย — คุณบังคับให้ network ส่ง message ครั้งเดียวเป๊ะภายใต้ crash และ retry ไม่ได้ สิ่งที่ สร้างได้จริง คือ idempotent consumer ที่ให้ผลลัพธ์ที่สังเกตได้เหมือนกันไม่ว่า event จะมาครั้งเดียวหรือร้อยครั้ง ซึ่งก็คือสิ่งที่
ApplyPaymentResultเป็นพอดี เวลาใครบอกว่าระบบของเขา “exactly-once” เกือบทุกครั้งเขาหมายถึงแบบนี้ — at-least-once delivery บวก idempotent handling - ไม่มี distributed rollback อัตโนมัติ นี่คือ choreography saga The Saga Handler → จึงไม่มี coordinator ที่ถือ logic การ compensate เมื่อ payment ล้มเหลว คำสั่งซื้อจะขยับ ไปข้างหน้า สู่
cancelledที่เป็น terminal state ไม่ใช่ ถอยหลัง ด้วยการ undo ขั้นตอนก่อนหน้า ในคอร์สนี้แค่นี้ก็พอ เพราะขั้นตอนก่อนหน้าเพียงอย่างเดียวที่มี side effect คือการสร้าง order row และcancelledก็เป็นจุดจบที่เหมาะสมอยู่แล้ว แต่ระบบที่ซับซ้อนกว่านี้ — ที่ reserve stock หรือ charge card จริงก่อนรู้ผล payment — จะต้อง compensate ด้วยการ release stock หรือ refund charge ซึ่งทั้งหมดเป็น action ไปข้างหน้า เพราะ refund คือ transaction ใหม่ ไม่ใช่การ undo ใน choreography แต่ละเซอร์วิสต้องตอบสนองต่อ failure event แล้ว compensate ขั้นตอนของตัวเอง โดยไม่มีจุดกลางคอยรับประกันว่าทุก compensation ได้รัน ส่วน orchestration saga จะยัด sequence นั้นไว้ใน state machine เดียว saga ตัวนี้จงใจไม่มี side effect ที่ต้อง compensate จึงยังไม่ต้องใช้ นี่คือขอบเขตที่ควรรู้ก่อนจะเพิ่มขั้นตอนที่ มี side effect เข้าไป
ข้อดีข้อเสีย
หัวข้อที่มีชื่อว่า “ข้อดีข้อเสีย”Idempotent consumer (dedupe + terminal-state guard) เทียบกับ เชื่อว่า delivery layer เป็น exactly-once
- Pros: ถูกต้องภายใต้การรับประกัน delivery ที่ระบบมีจริงคือ at-least-once แทนที่จะพึ่งการรับประกันที่ไม่มีทางมี ความปลอดภัยอยู่ใน database ภายใน transaction เดียว ซึ่งทนทานและตรวจสอบย้อนหลังได้ ไม่ใช่อยู่ในสมมติฐานเปราะ ๆ ว่า broker จะไม่ redeliver แถมยัง compose ต่อได้ เพราะ consumer ตัวไหนที่ dedupe บน
Event.IDก็ปลอดภัยทั้งจาก redelivery ของตัวเองและจาก producer ต้นทางที่ derive id แบบ deterministic - Cons: consumer ที่มี state ทุกตัวต้องพก table
processed_eventsมาด้วย พร้อมวินัยที่ต้องเช็คก่อนเสมอ ที่เป็น schema จริงและโค้ดจริงที่ต้องทำให้ถูก และ dedupe table จะโตไม่มีขอบเขตถ้าไม่มี retention policy การ pruneevent_idเก่าเป็นระยะอยู่นอกขอบเขตบทนี้ นี่คือต้นทุน operational เล็ก ๆ ที่จะไม่มีเลยถ้าเชื่อ broker แบบ exactly-once ในตำนาน สมมติว่า broker แบบนั้นมีอยู่จริง
Guard ทั้งสอง (dedupe และ terminal-state) เทียบกับ มีแค่ตัวเดียว
- Pros: สองตัวด้วยกันถูกต้องภายใต้ duplication และ out-of-order และ concurrent delivery — dedupe จัดการ same-id redelivery, การเช็ค terminal-state บวก
for updateจัดการ event คนละตัวที่แข่ง re-resolve คำสั่งซื้อที่จบไปแล้ว — ซึ่งคือชุดเต็มของสิ่งที่ at-least-once delivery ทำกับคุณได้จริง - Cons: มีชิ้นส่วนมากกว่าที่อ่านผ่านครั้งแรกแล้วรู้สึกว่าจำเป็น และ guard ทั้งสองดูซ้ำซ้อนจนกว่าจะลองสร้างเคสที่แต่ละตัวลำพังจะพลาด reviewer ที่เผลอลบ “ตัวที่ซ้ำซ้อน” ออกจะดึง bug ละเอียดอ่อนกลับเข้ามา ซึ่งโผล่เฉพาะภายใต้ interleaving บางแบบเท่านั้น และเป็น failure ประเภทที่จับได้ยากที่สุดใน testing
ติดตั้ง
หัวข้อที่มีชื่อว่า “ติดตั้ง”ไม่มีอะไรใหม่ต้องเขียน — guard สองตัวที่บทนี้พูดถึงอยู่ใน services/order/internal/repo/orders.go แล้วจาก The Saga Handler → และ processed_events ถูกสร้างโดย migrations/order/0002_processed_events.sql ของบทนั้น อ่าน guard ทั้งสองอีกครั้งด้วย framing ของบทนี้ คราวนี้เป็นภาพรวมแทนที่จะเป็นรายละเอียดที่ผ่านไป:
// Guard 1: dedupe. Already-seen event id -> commit the no-op, touch nothing.tag, err := tx.Exec(ctx, ` insert into processed_events (event_id) values ($1) on conflict do nothing`, eventID)if err != nil { return fmt.Errorf("repo: record processed event %s: %w", eventID, err)}if tag.RowsAffected() == 0 { return tx.Commit(ctx)}
// Guard 2: lock the order row, then refuse to re-resolve a terminal order.var status stringif err := tx.QueryRow(ctx, ` select status from orders where id = $1 for update`, orderID,).Scan(&status); err != nil { return fmt.Errorf("repo: lock order %s: %w", orderID, err)}if status != "pending" { return tx.Commit(ctx)}ถ้าก่อนหน้านี้คุณข้ามการ apply 0002_processed_events.sql ไป ก็แค่รัน:
create table processed_events ( event_id text primary key, processed_at timestamptz not null default now() );migrate -path migrations/order -database "$ORDER_DB_URL" upตรวจสอบผล
หัวข้อที่มีชื่อว่า “ตรวจสอบผล”ยก stack ขึ้นมาแล้วรัน Catalog, Order, gateway, และ Payment แบบเดียวกับที่ The Saga Handler → ทิ้งไว้:
cd deploy/compose && docker compose up -d postgres kafkago run ./services/catalog/cmdgo run ./services/order/cmdgo run ./gateway/cmdgo run ./services/payment/cmdสร้างคำสั่งซื้อขนาดเล็กแล้วปล่อยให้ loop resolve จนเป็น CONFIRMED แบบเดียวกับบทก่อน:
curl -s -X POST localhost:8080/v1/orders \ -H 'Content-Type: application/json' \ -d '{"customer_id":"cust-1","items":[{"product_id":"8f14e45f-ceea-4c9d-b2a5-0c1e3f4a9b21","quantity":2}]}'Poll สถานะจนได้ ORDER_STATUS_CONFIRMED แล้วจด id ของคำสั่งซื้อไว้ ทีนี้พิสูจน์ idempotency ตรง ๆ: redeliver event payment.succeeded เดิมเป๊ะ ลงบน topic payments ด้วย id เดิมที่ console consumer แสดงใน Process & Publish → (<order-id>:payment) produce record ดิบหนึ่งอันตรงเข้า Kafka:
docker compose exec -T kafka /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server localhost:9092 --topic payments <<'EOF'{"id":"3a7c9e21-1e4d-4b8a-9c6e-2f8b1d5a7c90:payment","type":"payment.succeeded","aggregate_id":"3a7c9e21-1e4d-4b8a-9c6e-2f8b1d5a7c90","payload":{"order_id":"3a7c9e21-1e4d-4b8a-9c6e-2f8b1d5a7c90","amount_cents":2598},"occurred_at":"2026-07-14T09:15:40Z"}EOFsaga จะ consume event นี้อีกครั้ง แต่เพราะ id อยู่ใน processed_events แล้ว Guard 1 จึง short-circuit และไม่มีอะไรเปลี่ยน ยืนยันว่า status ของคำสั่งซื้อไม่โดนแตะ และ event บันทึกไว้ครั้งเดียวพอดี:
docker compose exec -T postgres psql -U shopmicro -d orders -c \ "select status from orders where id = '3a7c9e21-1e4d-4b8a-9c6e-2f8b1d5a7c90';" status----------- confirmed(1 row)docker compose exec -T postgres psql -U shopmicro -d orders -c \ "select count(*) from processed_events where event_id = '3a7c9e21-1e4d-4b8a-9c6e-2f8b1d5a7c90:payment';" count------- 1(1 row)ยังเป็น confirmed, มี processed_events row เดียวพอดี, และไม่มี outbox row order.confirmed ที่สองถูกเขียน — การ redeliver เป็น no-op จริง ทีนี้ทดสอบ Guard 2 แทน Guard 1: produce event id ต่างกัน สำหรับคำสั่งซื้อเดิมที่ confirmed ไปแล้ว — payment.failed ที่มาสายที่ dedupe จะปล่อยผ่าน:
docker compose exec -T kafka /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server localhost:9092 --topic payments <<'EOF'{"id":"late-failure-001","type":"payment.failed","aggregate_id":"3a7c9e21-1e4d-4b8a-9c6e-2f8b1d5a7c90","payload":{"order_id":"3a7c9e21-1e4d-4b8a-9c6e-2f8b1d5a7c90","amount_cents":2598,"reason":"amount exceeds limit"},"occurred_at":"2026-07-14T09:20:00Z"}EOFid นี้เป็นของใหม่ Guard 1 จึงปล่อยผ่าน แต่คำสั่งซื้อเป็น confirmed ไปแล้ว if status != "pending" ของ Guard 2 จึงปฏิเสธที่จะขยับสถานะ เช็ค status อีกครั้ง:
docker compose exec -T postgres psql -U shopmicro -d orders -c \ "select status from orders where id = '3a7c9e21-1e4d-4b8a-9c6e-2f8b1d5a7c90';"ยังเป็น confirmed — event ที่มาสายและขัดแย้งพลิกคำสั่งซื้อที่ resolve แล้วไม่ได้ ซึ่งคือเคสที่ dedupe ลำพังจะพลาดพอดี แล้วยืนยันว่า module ยังคง build ผ่าน:
go build ./...ไม่มี output แปลว่าสำเร็จ
ตรวจสอบความเข้าใจ:
- ลาก trace
payment.succeededที่ redeliver (id เดิม) กับpayment.failedที่ใหม่จริง ๆ สำหรับคำสั่งซื้อที่confirmedแล้ว guard ตัวไหนหยุดแต่ละอัน และอะไรจะพังถ้า guard นั้นไม่อยู่? - ระบบเป็น “at-least-once delivery พร้อม effect แบบ exactly-once” ความต่างที่แม่นยำระหว่างสองวลีนั้นคืออะไร และทำไมคุณถึงมี exactly-once delivery แทนไม่ได้?
- ทำไม
for updateบนบรรทัดselect statusถึงจำเป็น แทนที่จะใช้แค่การเช็คif status != "pending"ลำพัง? ลองสร้าง interleaving ที่for updateป้องกันไว้ - saga นี้ไม่มี compensation logic และไม่มี orchestrator ต้องเพิ่มขั้นตอนใหม่แบบไหนเข้า order flow เรื่องนี้ถึงจะกลายเป็นปัญหาจริง และ compensation ของขั้นตอนนั้นจะหน้าตาเป็นอย่างไร?
OrderRepo.ApplyPaymentResult ปลอดภัยภายใต้ at-least-once delivery ของ Kafka เพราะ guard สองตัวใน transaction เดียว Guard 1, insert into processed_events ... on conflict do nothing, dedupe บน id ของ event การ redeliver event เดิม จะเจอ id ที่บันทึกไว้แล้ว commit แบบ no-op และไม่แตะอะไรทั้งสิ้น guard นี้ยังกัน redelivery ของ Payment เองได้ด้วย เพราะ Payment derive id ของ result event แต่ละตัวแบบ deterministic
Guard 2 คือ select status ... for update บวก if status != "pending" ซึ่งจับ row lock เพื่อ serialize concurrent consumer แล้วปฏิเสธการ re-resolve คำสั่งซื้อที่ terminal ไปแล้ว จึงดัก event คนละตัว ที่พยายามขยับคำสั่งซื้อซึ่ง dedupe จะปล่อยผ่าน เมื่อรวมสอง guard ไว้ใน pgx.Tx เดียวกับการ update status และ outbox row handler จึงถูกต้องภายใต้ duplication, reordering และ concurrency
สิ่งที่ระบบนี้จงใจ ไม่ เป็นคือ exactly-once delivery เพราะจริง ๆ แล้วเป็น at-least-once delivery พร้อม effect แบบ idempotent exactly-once ที่เป็นเวอร์ชันที่ทำได้จริงของเป้าหมายเดียวกัน และไม่มี distributed rollback อัตโนมัติ เพราะ choreography saga ที่ไม่มี side effect ให้ compensate จะขยับคำสั่งซื้อที่ล้มเหลว ไปข้างหน้า สู่ cancelled แทนการ undo ขั้นตอนก่อนหน้า การ redeliver event เดิมและ inject event ที่ขัดแย้งเข้ามาสาย พิสูจน์ทั้งสองอย่างพร้อมกัน — คำสั่งซื้อเปลี่ยน state ครั้งเดียวพอดี ไม่ว่า delivery layer จะทำอะไรก็ตาม นั่นทำให้ saga สมบูรณ์ บทถัดไป Notification Service → consume payments stream เดียวกันภายใต้ group ของตัวเองแล้วเปลี่ยนแต่ละผลลัพธ์เป็น notification ให้ลูกค้า — reader อิสระตัวที่สองของ event ที่ saga นี้ resolve อยู่มาตลอด