Consuming Orders
สิ่งที่จะสร้าง
หัวข้อที่มีชื่อว่า “สิ่งที่จะสร้าง”services/payment/cmd/main.go — composition root ของเซอร์วิส Payment รูปแบบเดียวกับที่ cmd/main.go ทุกตัวในคอร์สนี้ใช้มาตั้งแต่ The gRPC Server →: สร้าง dependency เริ่มงาน แล้ว shutdown อย่างนุ่มนวลบน SIGINT/SIGTERM แต่ Payment เป็น process คนละแบบกับที่เคยสร้างมาทั้งหมด เพราะไม่ใช่ gRPC server และไม่ได้ทำงานเพราะมี request เข้ามา trigger แต่เป็น kafka.Consumer ใน consumer group ของตัวเองชื่อ "payment" ที่อ่าน topic "orders" ซึ่ง Producer & Consumer → กับ Outbox & Relay → สร้างและต่อสายไว้แล้ว — งานทั้งหมดของ Payment เริ่มตั้งแต่วินาทีที่ event order.created มาถึงตรงนั้น
บทนี้ต่อสาย consumer, kafka.Publisher (สร้างไว้ตอนนี้ ใช้จริงบทถัดไป) และ closure handle แบบ placeholder ที่พิสูจน์ว่าการต่อสายใช้ได้ โดย filter หา event.Type == "order.created" แล้ว log ตัวที่เจอ ต่อมา Process & Publish → จะแทนที่ closure นั้นด้วย internal/processor.Processor.Handle ที่ตัดสินผลลัพธ์การชำระเงินจริงแล้ว publish ออกไป — แยกเป็น “ต่อสาย dependency ตอนนี้ เขียน logic จริงบทถัดไป” แบบเดียวกับที่ The gRPC Server → ใช้กับ Catalog client ของ Order
Payment เป็น event-driven และ stateless — ไม่มี database ไม่มี pkg/pg pool ที่ไหนในเซอร์วิสนี้เลย นี่คือการตัดสินใจออกแบบที่ตั้งใจและควรพูดให้ชัดตั้งแต่ต้น เพราะทุกเซอร์วิสก่อนหน้านี้ในคอร์ส (Catalog, Order) วนอยู่รอบ ๆ database และ gRPC คือทางเดียวที่ทุกอย่างคุยกัน Payment ไม่ได้ expose RPC สำหรับ “charge คำสั่งซื้อนี้” ที่ Order เรียกแบบ synchronous เหมือนที่ Order เรียก Catalog เพื่อขอราคา — แต่ outbox relay ของ Order publish order.created แล้ว Payment ซึ่ง subscribe topic เดียวกันภายใต้ consumer group ของตัวเอง ก็ตอบสนองตามตารางเวลาของตัวเอง ไม่มีเซอร์วิสไหนต้องรอให้อีกฝั่งออนไลน์อยู่ในวินาทีที่สร้างคำสั่งซื้อ Payment จะ redeploy, restart หรือ down ไปสักพักก็ได้ แล้ว order.created ทุกตัวที่พลาดไปก็ยังรออยู่ใน log ของ Kafka ตอนกลับมา เพราะ Topics, Partitions & Consumer Groups → พูดไว้แล้วว่าการ commit offset ไม่เคยลบ message ที่อยู่ข้างใต้
NewConsumer(brokers, "payment", "orders") คือคำตอบทั้งหมดของ “Payment เห็นคำสั่งซื้อทุกอันไหม ไม่ว่าเซอร์วิสอื่นจะทำอะไรกับ topic เดียวกันนี้” Topics, Partitions & Consumer Groups → ตั้งชื่อกฎนี้ไว้แล้วโดยยังไม่มี consumer ตัวที่สองมาพิสูจน์: consumer สองตัวใน consumer group เดียวกัน จะแบ่ง partition ของ topic กัน แต่ consumer สองตัวใน group ต่างกัน แต่ละตัวจะได้อ่าน stream ทั้งหมดแบบเป็นอิสระจากกันเอง groupID ของ Payment คือ "payment" — ไม่มีอะไรอื่นในระบบนี้จะใช้ชื่อนี้ร่วมด้วยเลย Payment จึงเห็น order.created ทุกตัวที่เคย publish มาอย่างครบถ้วน ไม่ว่า Notification (Module 9) หรือ saga ของ Order เอง (Module 10) จะ consume ไปแล้ว กำลังตามอยู่ หรือยังไม่มีตัวตนเลยก็ตาม นั่นคือการรับประกันที่ส่วน Verify ของบทนี้จะพิสูจน์ให้เห็น — Payment เห็น event ทันทีที่ publish โดยไม่ต้องประสานงานกับ Order เลย นอกจากตกลงชื่อ topic ให้ตรงกัน
อีกส่วนใหม่คือตัว filter เอง topicFor ของ Outbox & Relay → route event type ทุกตัวที่เป็น order.* — order.created, order.confirmed, order.cancelled — ไปยัง topic "orders" เดียวกัน consumer group ของ Payment เองก็รับทั้งหมด แต่ Payment มีงานต้องทำแค่ตอนคำสั่งซื้อถูกสร้างครั้งแรกเท่านั้น order.confirmed กับ order.cancelled เป็นเรื่องของคนอื่น if e.Type != "order.created" { return nil } คือการตัดสินใจทั้งหมดนั้น การ return nil แทน error ทำให้ offset ยัง commit ต่อ และ event นั้นจะไม่โดน redeliver อีก แค่ถูกเพิกเฉยอย่างตั้งใจเท่านั้น
ข้อดีข้อเสีย
หัวข้อที่มีชื่อว่า “ข้อดีข้อเสีย”Payment ตอบสนอง order.created ผ่าน Kafka เทียบกับ Order เรียก Payment แบบ synchronous ผ่าน gRPC (แบบเดียวกับที่ Order เรียก Catalog เพื่อขอราคา)
- Pros: Order ไม่ต้องพึ่ง Payment ให้เข้าถึงได้เพื่อจะวางคำสั่งซื้อเสร็จ — outage หรือ deploy ที่ช้าใน Payment ไม่กลายเป็น outage ใน Order แบบที่ The gRPC Server → ยอมรับความเสี่ยงนั้นชัด ๆ สำหรับ Catalog การเพิ่มเซอร์วิสที่สอง สาม หรือสิบที่ต้องตอบสนองต่อคำสั่งซื้อใหม่ทุกอัน (Notification, order saga) ไม่ทำให้ Order เสียอะไรเลย ไม่มีตัวไหนต้องแตะโค้ดของ Order แค่ subscribe topic เดียวกันก็พอ
- Cons: ไม่มี request/response — Order ไม่รู้จริง ๆ ว่า Payment อนุมัติหรือปฏิเสธการ charge ตอนที่
CreateOrderreturn แล้ว รู้ทีหลังตอนที่ event ผลลัพธ์ของ Payment เองมาถึง คอร์ส (หรือ product) ที่ต้องการให้ผู้เรียกเห็นผลการชำระเงินแบบ synchronous ใน HTTP response เดียวกัน จะต้องใช้ gRPC ตรงนี้แทน หรือใช้ client ที่ poll/subscribe ผลแยกต่างหาก
Filter event.Type ภายใน topic "orders" ที่ใช้ร่วมกัน เทียบกับ แยก topic ต่อ event type (order.created, order.confirmed, order.cancelled แต่ละตัวมี topic ของตัวเอง)
- Pros: topic เดียวเก็บ event ทุกตัวเกี่ยวกับคำสั่งซื้อไว้ใน log เดียวที่เรียงลำดับเข้มงวดต่อ partition — Outbox & Relay → พึ่งพาสิ่งนี้อยู่แล้วสำหรับการแบ่ง partition ด้วย
aggregate_idและ consumer ใหม่ที่อยากได้ประวัติเต็มของคำสั่งซื้อ (ทุกการเปลี่ยน status ตามลำดับ) อ่าน topic เดียวจบ ไม่ต้องอ่านสาม topic ยิ่งกว่านั้น การเพิ่ม event typeorder.*ใหม่ทีหลังก็ไม่ต้องจัดเตรียม topic ใหม่หรือต่อสาย consumer ใหม่ที่ไหนเลย เพราะ consumer ที่มีอยู่ได้รับ event นั้นอยู่แล้ว แค่เพิ่มcaseใหม่ หรือแบบในบทนี้คือปล่อยให้ filter เดิม ignore ต่อไป - Cons: consumer ทุกตัวของ
"orders"รวมถึง Payment ได้รับและต้อง deserialize แล้ว filter ทิ้ง event type ที่ไม่สนใจ — ต้นทุนคงที่เล็ก ๆ ต่อ message ที่การแยก topic ต่อ type จะหลีกเลี่ยงได้ เพราะ consumer จะ subscribe เฉพาะ topic ที่มี event ที่ต้องการเท่านั้น ในสเกลของระบบนี้ต้นทุนนั้นเล็กน้อยมาก แต่ระบบที่ throughput สูงกว่ามากและมี pattern การ consume ต่างกันชัดเจนต่อ event type การแยก topic อาจสมเหตุสมผลกว่า
ติดตั้ง
หัวข้อที่มีชื่อว่า “ติดตั้ง”1. services/payment/cmd/main.go
หัวข้อที่มีชื่อว่า “1. services/payment/cmd/main.go”// Command payment runs the Payment service: a Kafka consumer that reacts// to order events. This lesson wires the consumer and publisher; the next// lesson replaces the placeholder handler below with the real decision and// publish logic in internal/processor.package main
import ( "context" "log" "os" "os/signal" "strings" "syscall"
"github.com/avetavos/shopmicro/pkg/config" "github.com/avetavos/shopmicro/pkg/events" "github.com/avetavos/shopmicro/pkg/kafka")
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()
// Lesson-1 snapshot: internal/processor doesn't exist yet, so handle is // a placeholder that only proves the consumer is wired correctly — it // filters for "order.created" and logs it. Process & Publish replaces // this closure with processor.Processor.Handle, which decides the // order's payment outcome and publishes it through publisher. handle := func(ctx context.Context, e events.Event) error { if e.Type != "order.created" { return nil } log.Printf("payment: consumed order.created id=%s aggregate_id=%s", e.ID, e.AggregateID) return nil }
go func() { if err := consumer.Run(ctx, 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 consumer.Run block จนกว่า ctx จะโดน cancel จึงต้องรันใน goroutine ของตัวเอง เป็น pattern เดียวกับ “เริ่ม loop ที่รันยาวใน goroutine แล้วให้ cancel() ปลดตอน shutdown” ที่ go relay.Run(ctx) ของ Outbox & Relay → วางไว้แล้ว
เราสร้าง publisher พร้อม defer close ไว้ตรงนี้ ทั้งที่ handle ยังไม่เรียกใช้เลย ตรงกับที่ The gRPC Server → dial Catalog client ของ Order ล่วงหน้าเต็มหนึ่งบทก่อนที่ CreateOrder จะได้ใช้จริง — ต่อสาย dependency ไว้ตอนนี้ ส่วน logic ที่ต้องใช้ค่อยมาบทถัดไป
ตรวจสอบผล
หัวข้อที่มีชื่อว่า “ตรวจสอบผล”เปิด Postgres กับ Kafka แล้วรัน Catalog, Order, และ gateway ตามที่ REST Mapping → ทิ้งไว้ให้รันอยู่:
cd deploy/compose && docker compose up -d postgres kafkago run ./services/catalog/cmdgo run ./services/order/cmdgo run ./gateway/cmdเปิด terminal ที่สี่แล้วเริ่ม Payment:
go run ./services/payment/cmdpayment: consumer started, group=payment topic=ordersสร้าง product แล้ววางคำสั่งซื้อผ่าน REST ตามที่ Verify ของ REST Mapping → ทำไว้:
curl -s -X POST localhost:8080/v1/products \ -H 'Content-Type: application/json' \ -d '{"name":"Coffee Mug","description":"350ml ceramic mug","price_cents":1299,"stock":50}'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}]}'กลับไปที่ terminal ของ Payment ภายในไม่กี่วินาทีหลังคำสั่งซื้อถูกสร้าง:
payment: consumed order.created id=3a7c9e21-1e4d-4b8a-9c6e-2f8b1d5a7c90 aggregate_id=3a7c9e21-1e4d-4b8a-9c6e-2f8b1d5a7c90Payment เห็น event นี้ทันทีที่ outbox relay ของ Order publish ออกมา โดยไม่มีการเรียกจาก Order ไป Payment เกิดขึ้นเลยใน flow นี้ หยุด Payment ด้วย Ctrl-C แล้วดู log shutdown:
payment: shutting downแล้วยืนยันว่า module ยังคง build ผ่าน:
go build ./...ไม่มี output แปลว่าสำเร็จ
ตรวจสอบความเข้าใจ:
- ถ้า Notification (Module 9) เริ่ม
kafka.NewConsumer(brokers, "notification", "orders")ของตัวเองในวันพรุ่งนี้ Notification ต้องพึ่งความร่วมมือจาก Payment เพื่อจะเห็นorder.createdทุกตัวทั้งในอดีตและอนาคตไหม เพราะอะไร? - ทำไม
handleถึง returnnil— ไม่ใช่ error — สำหรับ eventorder.confirmedหรือorder.cancelledแทนที่จะเช่น log warning แล้วข้าม commit? - การรับประกันของ Payment จะเปลี่ยนไปอย่างไร ถ้าแชร์ consumer group (
"payment") กับ Payment instance ตัวที่สอง เทียบกับแชร์กับเซอร์วิสที่ต่างกันโดยสิ้นเชิง?
services/payment/cmd/main.go สร้าง kafka.Publisher, kafka.Consumer ใน consumer group ของตัวเองชื่อ "payment" ที่อ่าน topic "orders" แล้วรัน placeholder handler ใน goroutine shutdown อย่างสะอาดบน SIGINT/SIGTERM ผ่าน context.Context ที่ cancel ได้ — pattern เดียวกับที่ goroutine ระยะยาวทุกตัวในคอร์สนี้ใช้ filter ของ handler if e.Type != "order.created" { return nil } จำเป็นเพราะ Outbox & Relay → route event type ทุกตัวที่เป็น order.* ไป topic เดียวกันนี้ และ Payment สนใจแค่ตัวแรกเท่านั้น เพราะ "payment" เป็น consumer group ที่ไม่มีอะไรอื่นในระบบนี้ใช้ร่วม Payment จึงได้มุมมองที่สมบูรณ์และเป็นอิสระของคำสั่งซื้อทุกอันที่เคยถูกสร้าง — ยืนยันด้วยการสร้างคำสั่งซื้อผ่าน gateway แล้วดู log ของ Payment ตอบสนองภายในไม่กี่วินาที โดยไม่มีการเรียกตรงจาก Order ไป Payment ที่ไหนใน path เลย บทถัดไป Process & Publish → แทนที่ placeholder ด้วย internal/processor.Processor: การตัดสินใจแบบ deterministic ที่เปลี่ยน order.created แต่ละตัวให้เป็น event payment.succeeded หรือ payment.failed บน topic "payments"