ข้ามไปยังเนื้อหา

Process & Publish

services/payment/internal/processor/processor.goProcessor ที่ method Handle คือ logic จริงที่ closure placeholder ของ Consuming Orders → ทำหน้าที่แทนไว้: รับ event order.created มาแล้วตัดสินใจแบบ deterministic ว่าการชำระเงินของคำสั่งซื้อนั้นสำเร็จหรือล้มเหลว สร้าง events.Event ใหม่ที่บรรจุผลลัพธ์ แล้ว publish ไป Kafka topic "payments" จากนั้นเปลี่ยนแปลงเล็ก ๆ ที่ services/payment/cmd/main.go: สร้าง Processor แล้วต่อสาย consumer.Run(ctx, proc.Handle) แทน placeholder ของบทที่แล้ว

ไม่มีการเชื่อมต่อ gateway จริงตรงนี้ — ไม่มี Stripe ไม่มีการเรียก card network การตัดสินใจของ Payment เป็นกฎเดียวที่ตายตัว: อนุมัติทุกอย่างที่ไม่เกิน $5,000 ปฏิเสธทุกอย่างที่เกิน จงใจให้เรียบง่ายแบบนี้ และประเด็นของบทนี้ไม่ใช่ตัวกฎเอง — แต่เป็นสิ่งที่ตกผลึกออกมาจากการที่การตัดสินใจบริสุทธิ์ (pure — input เดิมให้ output เดิมเสมอ ไม่มี I/O ไม่มีความสุ่ม) ในระบบที่ event ทุกตัวอาจถูกส่งซ้ำได้มากกว่าหนึ่งครั้ง

Producer & Consumer → กำหนดไว้แล้วว่า Consumer.Run เป็น manual-commit แบบ at-least-once: ถ้า Payment crash หลัง Handle ทำงานเสร็จแต่ก่อน offset จะ commit Kafka จะ redeliver event order.created ตัวเดิมเป๊ะ ๆ ตอน restart consumer ทุกตัวที่สร้างบนการรับประกันนั้นต้อง idempotent เพราะการประมวลผล message ซ้ำต้องไม่ทำให้ลูกค้าโดน charge สองครั้ง หรือปล่อยผลลัพธ์ที่ขัดแย้งกันออกมาสองชุด

Processor.Handle ได้ idempotency มาฟรีโดยไม่ต้องเขียนโค้ดเพิ่มสักบรรทัด เพราะเป็น stateless และ deterministic อ่านแค่ event ที่เข้ามา ตัดสินใจจาก OrderCreated.TotalCents อย่างเดียว และไม่เขียนอะไรลง database เลย ต่อให้ redeliver order.created ตัวเดิมสิบครั้ง Handle ก็คำนวณ Result เดิมเป๊ะและ publish ผลลัพธ์เดิมเป๊ะทุกครั้ง ไม่มี state ตรงไหนให้ delivery ซ้ำทำเสียหาย เพราะไม่มี state เลย

รายละเอียดเดียวที่ทำให้เรื่องนี้ปลอดภัยสำหรับ consumer ปลายทางด้วย ไม่ใช่แค่ Payment เอง คือ ID: e.ID + ":payment" id ของ event payment.succeeded/payment.failed ใหม่ถูกสร้างมาจาก e.ID ของ event order.created ต้นทางเอง ไม่เคยสร้างใหม่ (ไม่มี uuid.New() ตรงนี้) นั่นหมายความว่า order.created ที่ถูก redeliver — e.ID เดิมทุกครั้ง — จะให้ event payment.* ที่มี id ที่สร้างมาเหมือนเดิมทุกครั้งด้วย consumer ปลายทางที่ dedupe บน events.Event.ID (pattern เดียวกับที่ Outbox & Relay → ตั้งไว้แล้วสำหรับการ redeliver ของ outbox relay เอง) จะจัดการ payment.succeeded ตัวที่สองสำหรับคำสั่งซื้อเดียวกันได้อย่างถูกต้องว่าเป็น no-op ที่จัดการไปแล้ว ไม่ใช่ผลการชำระเงินตัวที่สองที่ขัดแย้งกัน

payment gateway จริงทำแบบนี้ไม่ได้ การ charge card ไม่ใช่ operation ที่บริสุทธิ์ แต่เป็น side effect ภายนอกที่มี latency มี failure mode ของตัวเอง และมีเงินเคลื่อนไหวจริง เซอร์วิส Payment ระดับ production จึงต้องมี ตาราง payments ที่ persist คั่นอยู่ก่อนการเรียกนั้น: หนึ่งแถวต่อหนึ่งคำสั่งซื้อ unique constraint บน order_id (หรือ idempotency key ที่ส่งให้ gateway อย่างชัดเจน) ที่ตรวจสอบก่อนเรียกออกไป ตอน redeliver handler จะเจอแถวที่มีอยู่แล้ว เห็นว่าการชำระเงินถูกพยายามไปแล้ว แล้ว return ผลลัพธ์ที่เก็บไว้แทนที่จะ charge ซ้ำ ดีไซน์แบบ stateless ของบทนี้คือเวอร์ชันที่ง่ายที่สุดของ “ทำให้ redelivery ปลอดภัย” ซึ่งการเชื่อมต่อ gateway จริงจะเอา persistence-based dedupe มาต่อยอดทับ ไม่ใช่มาแทนที่ ประเด็นนี้พูดไว้ชัดเจนใน Pros & cons ด้านล่าง

การตัดสินใจแบบ stateless, deterministic (บทนี้) เทียบกับ ตาราง payments ที่ persist พร้อม dedupe (แนวทางจริงในโลกจริง)

  • Pros: infrastructure เป็นศูนย์ — ไม่มี database ไม่มี migration ไม่มี connection pool ไม่มีอะไรที่จะ disk เต็มหรือต้อง backup; idempotency เป็นผลลัพธ์ของดีไซน์แทนที่จะเป็นสิ่งที่ต้อง implement และ test แยกต่างหาก; ถูกต้องอย่างแท้จริงสำหรับกฎใดก็ตามที่ต้องดูแค่ event ที่เข้ามาเท่านั้น
  • Cons: ใช้กับอะไรที่มี side effect จริงไม่ได้ พอ “ตัดสินใจ” กลายเป็น “เรียก card network” การเรียกซ้ำตอน redeliver จะ charge ซ้ำจริง เว้นแต่จะมีอย่างอื่นคอยจำว่าเคยเกิดขึ้นแล้ว อีกทั้งยังเขียนกฎที่ต้องอาศัยข้อมูลมากกว่า event ตรงหน้าไม่ได้ — การตรวจ fraud จริงเทียบกับประวัติการชำระเงินของลูกค้า ยอดรวมที่วิ่งอยู่เทียบกับ spending limit หรือ “refund คำสั่งซื้อนี้ไปแล้วหรือยัง” ทั้งหมดต้องการ state ที่ Payment ไม่มีตรงนี้

การสร้าง id ของ event ผลลัพธ์จาก id ของ event ต้นทาง (e.ID + ":payment") เทียบกับ สร้าง id ใหม่ทุกครั้งที่ publish event (uuid.New())

  • Pros: id ที่สร้างมาเป็น deterministic ในตัวเอง — order.created เดิมให้ id payment.* เดิมเสมอ ทำให้ consumer ปลายทางทุกตัวได้ dedupe key ที่ถูกต้องฟรีเทียบกับการ redeliver ของ Payment เอง โดยไม่ต้องประสานงานระหว่าง Payment กับ consumer ปลายทางเลย นอกจากข้อตกลงเดียวว่า “dedupe บน Event.ID” ที่เป็นกฎที่ Outbox & Relay → สอนไว้แล้ว
  • Cons: id ใหม่ต่อการ publish จะดู “ปกติ” กว่าสำหรับ event ที่โดยแนวคิดแล้วเป็นข้อเท็จจริงใหม่ (“payment succeeded” อาจนับเป็น event ของตัวเอง ไม่ใช่การพูดซ้ำของ “order created”) และจะไม่ชนกันเลยถ้าคำสั่งซื้อสองอันที่ต่างกันจริง ๆ ให้ string ที่สร้างมาเหมือนกัน — ความเสี่ยงที่จริง ๆ แล้วเป็นศูนย์ตรงนี้เพราะ e.ID unique อยู่แล้วทั่วทั้งระบบ แต่ก็ควรพูดถึงในฐานะ trade-off ทั่วไประหว่าง id ที่สร้างมาจากของเดิมกับ id ที่สร้างใหม่สด ๆ
// Package processor implements Payment's business logic: given an
// order.created event, decide deterministically whether that order's
// payment succeeds or fails, and publish the result to Kafka's "payments"
// topic.
package processor
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/avetavos/shopmicro/pkg/events"
"github.com/avetavos/shopmicro/pkg/kafka"
)
// limitCents is the largest order Payment approves. Orders over $5,000
// fail with Reason "amount exceeds limit".
const limitCents = 500000
// OrderCreated is the payload of an "order.created" event, as published by
// Order's outbox relay.
type OrderCreated struct {
OrderID string `json:"order_id"`
CustomerID string `json:"customer_id"`
TotalCents int64 `json:"total_cents"`
}
// Result is the payload of a "payment.succeeded" or "payment.failed"
// event.
type Result struct {
OrderID string `json:"order_id"`
AmountCents int64 `json:"amount_cents"`
Reason string `json:"reason,omitempty"`
}
// Processor decides each order's payment outcome and publishes it.
type Processor struct {
pub *kafka.Publisher
}
// New returns a Processor that publishes results through pub.
func New(pub *kafka.Publisher) *Processor {
return &Processor{pub: pub}
}
// Handle implements the Consumer.Run handle signature. It ignores every
// event.Type except "order.created", decides that order's payment outcome
// deterministically from TotalCents alone, and publishes a
// "payment.succeeded" or "payment.failed" event to the "payments" topic,
// keyed by OrderID.
func (p *Processor) Handle(ctx context.Context, e events.Event) error {
if e.Type != "order.created" {
return nil
}
var oc OrderCreated
if err := json.Unmarshal(e.Payload, &oc); err != nil {
return fmt.Errorf("processor: unmarshal order.created payload: %w", err)
}
eventType := "payment.succeeded"
result := Result{OrderID: oc.OrderID, AmountCents: oc.TotalCents}
if oc.TotalCents > limitCents {
eventType = "payment.failed"
result.Reason = "amount exceeds limit"
}
payload, err := json.Marshal(result)
if err != nil {
return fmt.Errorf("processor: marshal result: %w", err)
}
out := events.Event{
ID: e.ID + ":payment",
Type: eventType,
AggregateID: oc.OrderID,
Payload: payload,
OccurredAt: time.Now().UTC(),
}
value, err := json.Marshal(out)
if err != nil {
return fmt.Errorf("processor: marshal event: %w", err)
}
if err := p.pub.Publish(ctx, "payments", oc.OrderID, value); err != nil {
return fmt.Errorf("processor: publish %s: %w", eventType, err)
}
return nil
}

บันทึกเป็น services/payment/internal/processor/processor.go p.pub.Publish(ctx, "payments", oc.OrderID, value) ใช้ oc.OrderID เป็น key ของ message — aggregate_id เดียวกับที่ Outbox & Relay → ใช้ key event order.* — ดังนั้นการรับประกันการเรียงลำดับต่อ partition ของ Topics, Partitions & Consumer Groups → จึงใช้ได้กับผลการชำระเงินของคำสั่งซื้อ เหมือนที่ใช้ได้กับ event ของคำสั่งซื้อใบเดียวกัน

// Command payment runs the Payment service: a Kafka consumer that reacts
// to order events, decides each order's payment outcome, and publishes the
// result back to Kafka.
package main
import (
"context"
"log"
"os"
"os/signal"
"strings"
"syscall"
"github.com/avetavos/shopmicro/pkg/config"
"github.com/avetavos/shopmicro/pkg/kafka"
"github.com/avetavos/shopmicro/services/payment/internal/processor"
)
func main() {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
brokers := strings.Split(config.Get("KAFKA_BROKERS", "localhost:9092"), ",")
publisher := kafka.NewPublisher(brokers)
defer publisher.Close()
consumer := kafka.NewConsumer(brokers, "payment", "orders")
defer consumer.Close()
proc := processor.New(publisher)
go func() {
if err := consumer.Run(ctx, proc.Handle); err != nil {
log.Printf("payment: consumer stopped: %v", err)
}
}()
log.Println("payment: consumer started, group=payment topic=orders")
stop := make(chan os.Signal, 1)
signal.Notify(stop, syscall.SIGINT, syscall.SIGTERM)
<-stop
log.Println("payment: shutting down")
cancel()
}

บันทึกทับ services/payment/cmd/main.go การเปลี่ยนแปลงเดียวจาก Consuming Orders →: proc := processor.New(publisher) แทนที่ closure handle แบบ placeholder และ consumer.Run(ctx, proc.Handle) ส่ง Processor.Handle เข้าไปตรง ๆ เพราะ signature เป็น func(context.Context, events.Event) error ซึ่งตรงกับที่ Consumer.Run ต้องการอยู่แล้ว จึงไม่มีอะไรในการต่อสายที่ต้องเปลี่ยนอีก

เปิด Postgres กับ Kafka แล้วรัน Catalog, Order, gateway, และ Payment:

Terminal window
cd deploy/compose && docker compose up -d postgres kafka
Terminal window
go run ./services/catalog/cmd
Terminal window
go run ./services/order/cmd
Terminal window
go run ./gateway/cmd
Terminal window
go run ./services/payment/cmd

เปิด terminal ที่ห้าแล้วดู topic "payments" โดยตรง:

Terminal window
docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 --topic payments --from-beginning

สร้างคำสั่งซื้อเล็ก ๆ — flow ของ Coffee Mug เดิมจาก Consuming Orders → ซึ่งไม่ถึง $5,000:

Terminal window
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}]}'

console consumer จะพิมพ์ผลลัพธ์ภายในไม่กี่วินาที:

{"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"}

จากนั้นสร้าง product ราคาแพงแล้วสั่งซื้อจำนวนที่ทำให้เกินลิมิต:

Terminal window
curl -s -X POST localhost:8080/v1/products \
-H 'Content-Type: application/json' \
-d '{"name":"Server Rack","description":"42U enterprise rack","price_cents":600000,"stock":5}'
Terminal window
curl -s -X POST localhost:8080/v1/orders \
-H 'Content-Type: application/json' \
-d '{"customer_id":"cust-1","items":[{"product_id":"7b23f5a1-4c9d-4e8a-b2a5-1e3f4a9b21c8","quantity":1}]}'
{"id":"9e4a2c11-8b3d-4f7a-a1c6-3d8b1e5a7c40:payment","type":"payment.failed","aggregate_id":"9e4a2c11-8b3d-4f7a-a1c6-3d8b1e5a7c40","payload":{"order_id":"9e4a2c11-8b3d-4f7a-a1c6-3d8b1e5a7c40","amount_cents":600000,"reason":"amount exceeds limit"},"occurred_at":"2026-07-14T09:16:12Z"}

600000 cents คือ $6,000 — เกิน threshold limitCents — ดังนั้น Handle จึง publish payment.failed พร้อม reason ตรงข้ามกับคำสั่งซื้อ Coffee Mug ที่เล็กกว่าซึ่ง publish payment.succeeded โดยไม่มี reason หยุด console consumer กับ Payment ด้วย Ctrl-C ในแต่ละ terminal

แล้วยืนยันว่า module ยังคง build ผ่าน:

Terminal window
go build ./...

ไม่มี output แปลว่าสำเร็จ

ตรวจสอบความเข้าใจ:

  • ถ้า Payment crash ทันทีหลัง publish payment.succeeded แต่ก่อนจะ commit offset ของ order.created poll รอบถัดไปจะ redeliver อะไร แล้ว id ของ event payment.succeeded ตัวที่สองจะหน้าตาเป็นอย่างไรเทียบกับตัวแรก?
  • ทำไมการสร้าง ID ของ event ผลลัพธ์จาก ID ของ event ต้นทาง (e.ID + ":payment") ปลอดภัยกว่าสำหรับ idempotency ปลายทาง เทียบกับการเรียก id generator ใหม่ข้างใน Handle?
  • การเชื่อมต่อ payment gateway จริงเป็น pure แบบที่ Handle ของบทนี้เป็นไม่ได้ เวอร์ชัน production ต้องเพิ่มอะไรอย่างน้อยที่สุดเพื่อให้ปลอดภัยภายใต้ at-least-once redelivery เดียวกัน?

Processor.Handle ของ services/payment/internal/processor/processor.go เพิกเฉยต่อ event ทุกตัวยกเว้น order.created ตัดสินใจผลลัพธ์ของคำสั่งซื้อจาก TotalCents เพียงอย่างเดียว — payment.succeeded ถ้าไม่เกิน $5,000, payment.failed พร้อม reason: "amount exceeds limit" ถ้าเกิน — แล้ว publish ผลลัพธ์ไป Kafka topic "payments" โดย key ด้วย OrderID เพื่อการรับประกันการเรียงลำดับต่อคำสั่งซื้อเดียวกันกับที่ "orders" พึ่งพาอยู่แล้ว services/payment/cmd/main.go ตอนนี้ต่อสาย processor.New(publisher) แล้วส่ง proc.Handle ตรงไปที่ consumer.Run แทนที่ placeholder ของ Consuming Orders → เพราะ Handle เป็น stateless และ deterministic จึง idempotent มาฟรีภายใต้ at-least-once redelivery ของ Consumer.Run ไม่ต้องมี persistence ไม่ต้องมีตาราง dedupe แค่ input เดิมให้ output เดิมเสมอ รวมถึง event ID ที่ derive ออกมาเหมือนเดิมทุกครั้งเพื่อให้ consumer ปลายทาง dedupe ได้

การเชื่อมต่อ gateway จริงจะต้องแลกความเรียบง่ายนี้กับตาราง payments ที่ persist และ dedupe key ที่ชัดเจน เพราะการ charge card เป็น side effect ที่ดีไซน์ของบทนี้ไม่เคยตั้งใจจะทำให้ปลอดภัยด้วยตัวเอง เท่านี้ก็ปิด Module 8 — Payment เปลี่ยน order.created ทุกตัวให้เป็นข้อเท็จจริง payment.succeeded หรือ payment.failed บน topic "payments" ได้อย่างน่าเชื่อถือ บทถัดไป Notification → (Module 9) consume topic เดียวกันนั้นแล้วเข้าคิวงานส่งบน RabbitMQ แทนที่จะส่งอีเมลตรง ๆ และ Order Saga → (Module 10) ก็ consume topic เดียวกันนี้ เพื่อเปลี่ยนสถานะคำสั่งซื้อแต่ละใบเป็น CONFIRMED หรือ CANCELLED