The Send Worker
สิ่งที่จะสร้าง
หัวข้อที่มีชื่อว่า “สิ่งที่จะสร้าง”services/notification/cmd/worker/main.go และ services/notification/internal/sender/sender.go — Notification worker, binary แยกตัวที่สองในเซอร์วิสเดียวกัน ตรงที่ process ของ Consuming events → ผลิต job ลงบน notification.send binary ตัวนี้ทำหน้าที่ consume โดยเรียก amqp.Client.Consume("notification.send", ...) — method ที่ Exchanges & Queues → สัญญาไว้ว่า “เซอร์วิส Notification จะใช้แบบไม่แก้” — แล้วส่ง delivery แต่ละตัวให้ Sender ที่ unmarshal SendJob ออกมา “ส่ง” ออกไป จากนั้นปล่อยให้ logic ack-on-success / nack-to-dead-letter ที่ Consume มีอยู่แล้วจัดการที่เหลือ
“ส่ง” ตรงนี้คือ log บรรทัดเดียว ไม่ใช่ SMTP หรือ SMS จริง — ตัวแทนที่ตรงไปตรงมาเดียวกับที่ Process & Publish → ใช้ตอนตัดสินการชำระเงินด้วยกฎ TotalCents แทน card network จริง ประเด็นของบทนี้ไม่ใช่การต่อ SendGrid แต่คือ รูปร่าง: process แยกที่ดึง work queue ที่ crash, scale, retry, และ dead-letter ได้อย่างเป็นอิสระจากฝั่ง Kafka โดยสิ้นเชิง
เซอร์วิส Notification กับ Notification worker เป็น สอง process อย่างตั้งใจ สถาปัตยกรรมวาดไว้เป็นสองกล่องตั้งแต่แรก Architecture → และการแยกนี้คือเหตุผลทั้งหมดที่เรา enqueue job แทนการส่ง inline
ฝั่ง Kafka (Consuming events →) ต้องเร็ว เพราะ commit offset หลัง handler return เท่านั้น อะไรที่ช้าอยู่ใน handler จึงหยุดการ consume "payments" ทั้งสาย ส่วนฝั่งการส่งตรงกันข้าม — ช้าและพังได้โดยธรรมชาติ เพราะต้องคุยกับ provider ภายนอกที่อาจ down โดน rate-limit หรือแค่ latency สูง
การแยกเป็นคนละ process ทำให้ของที่ช้า (การส่ง) ไม่มีวันหยุดของที่ต้องเร็ว (การ consume Kafka) และแต่ละฝั่ง scale บนแกนของตัวเอง: เพิ่ม Kafka partition เพื่อ throughput ของ event เพิ่ม worker เพื่อ throughput ของการส่ง
แกนที่สองนี้คือสิ่งที่ work queue ของ RabbitMQ มีไว้เพื่อ Exchanges & Queues → สร้าง notification.send เป็น queue แบบ competing-consumers รัน worker ตัวเดียวก็ได้ job ไปทั้งหมด รันสามตัว RabbitMQ ก็แบ่ง job ให้แต่ละตัวทำราวหนึ่งในสาม โดยมี prefetchCount เป็น 10 ที่ Acks, Retry & Dead Letters → set ไว้คอยคุมให้การแบ่งยุติธรรม
การเพิ่ม worker เพิ่มความจุการส่งได้โดย ไม่ ต้องแก้โค้ดและ ไม่ ต้องประสานงานอะไร แค่เริ่ม binary ตัวนี้อีกก๊อปเท่านั้น นี่คือด้านตรงข้ามกับ Kafka อย่างตั้งใจ เพราะ parallelism ของ consumer group ติดเพดานที่จำนวน partition ของ topic ส่วน work queue ของ RabbitMQ scale ตามจำนวน process ล้วน ๆ ซึ่งเข้ากับโจทย์ “มีกอง job รออยู่ กระจายไปบน worker เท่าที่มี” พอดี
worker สืบทอดความถูกต้องมาจาก Consume ที่ไม่ได้แก้อะไร และควรพูดให้ชัดว่าได้อะไรมา — พร้อมกับต้องแลกด้วยอะไร Consume ใช้ autoAck=false job จึงหลุดออกจาก queue ก็ต่อเมื่อ Sender.Deliver return nil ถ้า worker crash กลางการส่ง RabbitMQ จะ redeliver job นั้นให้ worker อีกตัว Acks, Retry & Dead Letters → นั่นคือ at-least-once ที่เป็นการรับประกันเดียวกับที่ Kafka ให้ขาเข้า
การที่ทั้งสอง hop เป็น at-least-once คือจุดสะดุด ในกรณีเลวร้ายสุด คำสั่งซื้อเดียวผลิต payment.succeeded ซ้ำ (Kafka redeliver) → send-job ซ้ำ (notifier รันสองครั้ง) → การส่งซ้ำ (worker รันสองครั้ง) Deliver ระดับ production จึงต้อง dedupe ด้วย table sent_notifications ที่ key ด้วย order_id + outcome แล้วเช็คก่อนส่ง แบบเดียวกับ persisted idempotency ที่ Process & Publish → บอกว่าการ charge card จริงต้องมี
Deliver ที่แค่ log ในบทนี้ idempotent อยู่แล้วโดยธรรมชาติ เพราะ log สองครั้งไม่เสียหาย เราจึงยังไม่สร้าง table นั้น แต่นี่คือสิ่งแรกที่ต้องเพิ่มทันทีที่ต่อกับ provider จริง
failure path ก็ต่อสายไว้ให้แล้ว ถ้า Deliver return error Consume จะเรียก Nack(false, false) ซึ่ง route job ไป notification.send.dead ผ่าน dead-letter argument ที่ DeclareTopology set ไว้ ไม่มี requeue loop ไม่มี job หาย มีแค่ job ที่ไปนอนรออยู่ในที่ที่คนเข้าไปดูได้ Acks, Retry & Dead Letters →
สำหรับ worker แล้ว failure แบบ “job นี้ไม่มีวันสำเร็จ” ตามธรรมชาติคือ body ที่ไม่ใช่ SendJob ที่ถูกต้อง หรือที่เรียกว่า poison message การโยนเข้า dead-letter ตั้งแต่ครั้งแรกจึงเป็นทางที่ถูกต้อง ส่วน failure แบบ transient อย่าง provider timeout ครั้งเดียว คือเคสที่ dead-letter ทันทีจัดการได้ไม่ดี และเป็นจุดที่ pattern TTL+DLX-cycling หรือ retry-count-header จาก Module 7 จะเสียบเข้ามา เรายังไม่ได้ต่อสายไว้ที่นี่ และ Pros & cons ด้านล่างพูดถึงช่องว่างนั้นตรง ๆ
ข้อดีข้อเสีย
หัวข้อที่มีชื่อว่า “ข้อดีข้อเสีย”Worker binary แยกต่างหาก เทียบกับ ส่ง inline ใน Kafka consumer ของ Notification
- Pros: ขั้นตอนการส่งที่ช้าและพังได้รันใน process ของตัวเอง ดังนั้น provider ที่ค้างหรือถูก rate-limit ไม่มีวันหยุดการ consume
"payments"ของ Kafka หรือ freeze offset progress; การส่ง scale เป็นอิสระแบบ competing consumers — เพิ่ม worker ไม่ใช่ partition ตอนมี backlog ต้องเคลียร์; และ crash กลางการส่งก็ไม่ทำอะไรหาย เพราะ job ยังไม่ ack บนnotification.sendRabbitMQ จึง redeliver ให้ใหม่ - Cons: ต้องมี binary ตัวที่สองให้ build, deploy และดูแล ทั้งที่เป็นฟีเจอร์เชิงตรรกะเดียว และ JSON ของ
SendJobกลายเป็น contract ระหว่างสอง process ไม่ใช่ระหว่างสอง function ใน process เดียวอีกต่อไป เปลี่ยนรูปร่างเมื่อไรก็ต้อง roll ทั้ง producer และ worker ทุกตัวพร้อมกัน หรือไม่ก็ทำ version ให้ message
Dead-letter job ที่ล้มเหลวตั้งแต่ครั้งแรก (Nack(false, false)) เทียบกับ retry แบบหน่วงเวลาอัตโนมัติ (TTL+DLX cycling หรือ retry-count header)
- Pros: ง่ายสุด ๆ และถูกต้องพอดีสำหรับ poison message — body ที่ไม่มีวัน parse หรือ job ที่จะถูกปฏิเสธเสมอจะออกจาก live queue หลังหนึ่งครั้งแล้วไปตกที่
notification.send.deadเห็นได้ ตรวจสอบได้ ไม่เคยบล็อก job ที่อยู่ข้างหลัง Acks, Retry & Dead Letters → - Cons: failure แบบ transient ล้วน ๆ — email provider timeout ครั้งเดียวและน่าจะสำเร็จถ้า retry อีกไม่กี่วินาที — ก็โดน dead-letter ตั้งแต่ครั้งแรกเหมือนกัน ไม่มีการลองใหม่อัตโนมัติ การกู้คืนจึงต้องรอให้มีคนไปเห็น job นั้นใน dead-letter queue แล้ว republish เอง queue ที่ downstream flaky พอจะอยากได้กลไก TTL+DLX หรือ retry-count จาก Module 7 มาซ้อนก่อน;
notification.sendยังไม่ต้องการ
ติดตั้ง
หัวข้อที่มีชื่อว่า “ติดตั้ง”1. services/notification/internal/sender/sender.go
หัวข้อที่มีชื่อว่า “1. services/notification/internal/sender/sender.go”// Package sender is the Notification worker's core: it takes one raw// notification.send job body, parses it back into a SendJob, and// "delivers" the notification. Delivery here is a log line standing in for// a real email/SMS provider — the mechanics (parse, deliver, succeed or// fail) are identical to a real integration; only the last step differs.package sender
import ( "context" "encoding/json" "fmt" "log"
"github.com/avetavos/shopmicro/services/notification/internal/notifier")
// Sender delivers notification.send jobs.type Sender struct{}
// New returns a Sender.func New() *Sender { return &Sender{}}
// Deliver implements the amqp.Client.Consume handle signature. It parses// the raw job body into a notifier.SendJob and delivers it. Returning nil// tells Consume to ack the message (it's removed from the queue);// returning an error tells Consume to Nack(requeue=false), routing the job// to notification.send.dead. A body that isn't a valid SendJob is a poison// message — it will never parse no matter how many times it's retried, so// dead-lettering it on the first failure is exactly the right call.func (s *Sender) Deliver(_ context.Context, body []byte) error { var job notifier.SendJob if err := json.Unmarshal(body, &job); err != nil { return fmt.Errorf("sender: unmarshal send-job: %w", err) } if job.OrderID == "" { return fmt.Errorf("sender: send-job missing order_id") }
// In a real service this is where an email/SMS provider gets called, // and where a persisted dedupe (keyed by order_id + outcome) would // guard against redelivery sending the same notification twice. log.Printf("sender: delivered notification for order %s (%s): %s", job.OrderID, job.Outcome, job.Message)
return nil}บันทึกไฟล์นี้เป็น services/notification/internal/sender/sender.go ไฟล์นี้ import notifier.SendJob เข้ามาใช้ เพราะ struct ของ producer คือ single source of truth ของรูปร่าง job worker กับเซอร์วิส Notification จึงไม่มีวัน drift ออกจากกัน
signature ของ Deliver คือ func(context.Context, []byte) error ซึ่งตรงกับที่ Client.Consume เรียกด้วย delivery แต่ละตัวพอดี เป็นการแยก “handler บาง ๆ แล้ว delegate งานจริง” แบบเดียวกับที่ consumer ทุกตัวในคอร์สนี้ใช้ เพียงแต่ลงลึกไปอีกหนึ่ง broker
2. services/notification/cmd/worker/main.go
หัวข้อที่มีชื่อว่า “2. services/notification/cmd/worker/main.go”// Command worker runs the Notification worker: it drains the// notification.send RabbitMQ queue and delivers each job. It's a separate// binary from the Notification service (cmd) on purpose — delivery is slow// and failable and must scale and fail independently of the Kafka consumer// that fills the queue. Run as many copies as you need; RabbitMQ splits the// jobs across them as competing consumers.package main
import ( "context" "log" "os" "os/signal" "syscall"
"github.com/avetavos/shopmicro/pkg/amqp" "github.com/avetavos/shopmicro/pkg/config" "github.com/avetavos/shopmicro/services/notification/internal/sender")
func main() { rabbitURL := config.Get("RABBITMQ_URL", "amqp://shopmicro:shopmicro@localhost:5672/")
client, err := amqp.Connect(rabbitURL) if err != nil { log.Fatalf("worker: connect to rabbitmq: %v", err) } defer client.Close()
if err := client.DeclareTopology(); err != nil { log.Fatalf("worker: declare rabbitmq topology: %v", err) }
s := sender.New()
go func() { if err := client.Consume("notification.send", s.Deliver); err != nil { log.Printf("worker: consume stopped: %v", err) } }()
log.Println("worker: consuming notification.send")
stop := make(chan os.Signal, 1) signal.Notify(stop, syscall.SIGINT, syscall.SIGTERM) <-stop
log.Println("worker: shutting down")}บันทึกไฟล์นี้เป็น services/notification/cmd/worker/main.go สังเกตว่าไม่มี kafka import อยู่ที่ไหนใน worker เลย เพราะแตะแค่ RabbitMQ อย่างเดียว ซึ่งคือประเด็นทั้งหมดของการแยกออกมา
client.Consume block จนกว่า delivery channel จะปิด จึงต้องรันใน goroutine เมื่อได้ SIGINT/SIGTERM main จะตกลงไปที่ client.Close() ที่ defer ไว้ ซึ่งปิด channel แล้วปลด block ของ Consume
worker เรียก DeclareTopology() ด้วย และเพราะเป็น idempotent จึงไม่มีผลข้างเคียง ไม่ว่าเซอร์วิส Notification จะ declare queue ไว้ก่อนแล้ว หรือ worker จะเป็นตัวแรกที่เริ่ม ไม่ว่าทางไหน notification.send และ notification.send.dead ก็มีอยู่ก่อนที่ worker จะเริ่ม consume เสมอ
ตรวจสอบผล
หัวข้อที่มีชื่อว่า “ตรวจสอบผล”ให้ทั้ง pipeline จาก Consuming events → รันอยู่ — Postgres, Kafka, RabbitMQ, และ Catalog, Order, gateway, Payment, และเซอร์วิส Notification ถ้าคุณทิ้ง job ไว้บน notification.send จากบทนั้นสักหนึ่งสองอันก็ยิ่งดี; ถ้าไม่ วางคำสั่งซื้อขนาดเล็กตอนนี้เพื่อ enqueue หนึ่งอัน:
cd deploy/compose && docker compose up -d postgres kafka rabbitmqgo run ./services/notification/cmdcurl -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}]}'ทีนี้เริ่ม worker อะไรก็ตามที่รออยู่บน notification.send จะถูกดึงออกทันที:
go run ./services/notification/cmd/workerworker: consuming notification.sendsender: delivered notification for order 3a7c9e21-1e4d-4b8a-9c6e-2f8b1d5a7c90 (succeeded): Your order 3a7c9e21-1e4d-4b8a-9c6e-2f8b1d5a7c90 is confirmed — payment of 2598 cents went through.นั่นคือทั้งระบบ end to end ทำงานเอง: REST POST /v1/orders กลายเป็น outbox row order.created → Kafka orders → Payment → Kafka payments → เซอร์วิส Notification → job notification.send ของ RabbitMQ → worker ตัวนี้ส่งออกไป รวมห้าเซอร์วิสและสอง broker โดย client ไม่ต้องรออะไรเลยสักขั้น ใน management UI ที่ http://localhost:15672 → Queues and Streams notification.send จะกลับมาที่ 0 messages ready เพราะ worker ack แล้วและ RabbitMQ ลบ job ทิ้ง
ทีนี้พิสูจน์ property competing-consumers หยุด worker แล้วเปิด terminal ใหม่ สอง อัน แต่ละอันรัน worker หนึ่งตัว:
go run ./services/notification/cmd/workergo run ./services/notification/cmd/workerยิง burst ของคำสั่งซื้อเพื่อให้หลาย job ตกบน queue พร้อมกัน:
for i in 1 2 3 4 5 6; do 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":1}]}' > /dev/nulldoneดู terminal ของ worker สองตัว: บรรทัด sender: delivered ... หกบรรทัด ถูกแบ่งข้ามทั้งสอง process ไม่ใช่ซ้ำในแต่ละตัว — RabbitMQ ส่งแต่ละ job ให้ worker ตัวเดียวเท่านั้น การ load-balance ที่ model ของ Kafka ให้คุณไม่ได้บน partition เดียว หยุด worker ทั้งสองด้วย Ctrl-C:
worker: shutting downdead-letter path คือกลไก Consume เดียวกับที่ Acks, Retry & Dead Letters → พิสูจน์ภายใต้ failure ไปแล้ว: job ที่ body ไม่ใช่ SendJob ที่ถูกต้อง หรือขาด order_id จะทำให้ Deliver return error แล้ว Consume ก็ nack โดยไม่ requeue job จึงไปตกที่ notification.send.dead ให้ตามตรวจสอบ เหมือนกับ demonstration -fail ใน Module 7 แต่ขยับขึ้นมาอีกหนึ่งชั้น คุณยืนยันได้ว่ามีทั้ง queue และ dead-letter queue คู่กันใน management UI ใต้ Queues and Streams
แล้วยืนยันว่า module ยังคง build ผ่าน:
go build ./...ไม่มี output แปลว่าสำเร็จ
ตรวจสอบความเข้าใจ:
- คุณรัน worker สองตัวและ burst ของคำสั่งซื้อหกอันผลิตการส่งรวมหกครั้ง ไม่ใช่สิบสอง นั่นคือ property อะไรของ RabbitMQ และทำไม single-partition Kafka topic ให้ load-split เดียวกันข้าม consumer สองตัวใน group เดียวไม่ได้?
- ทั้ง Kafka hop และ RabbitMQ hop เป็น at-least-once ลาก trace กรณีเลวร้ายสุดที่คำสั่งซื้อเดียวส่งผลให้ notification เดียวกันถูก ส่งสองครั้ง
Deliverระดับ production จะเพิ่มอะไรเพื่อให้ปลอดภัย? - ทำไม body ที่ล้มเหลว
json.Unmarshalถึงถูก dead-letter ตั้งแต่ครั้งแรกแทนที่จะ retry ส่วน failure แบบ “email provider timeout” จริง ๆ ก็ไม่ควรถูกทำแบบนั้น? - worker import
notifier.SendJobแทนที่จะ declare struct ก๊อปของตัวเอง อะไรจะพังเป็นอย่างแรกถ้า producer กับ worker เห็นรูปร่างไม่ตรงกัน และการ import ร่วมกันช่วยกันไว้อย่างไร?
services/notification/cmd/worker/main.go คือ binary แยกตัวที่ไม่มี Kafka อยู่ในนั้นเลย หน้าที่คือเรียก amqp.Client.Consume("notification.send", sender.Deliver) ที่เป็น method pkg/amqp ตัวเดียวกับที่ Exchanges & Queues → สร้างไว้และสัญญาว่า Notification จะใช้แบบไม่แก้
Sender.Deliver ใน services/notification/internal/sender/sender.go unmarshal body ดิบแต่ละอันกลับเป็น notifier.SendJob ของ producer เอง “ส่ง” ออกไปด้วย log บรรทัดเดียวแทน provider จริง แล้ว return nil เพื่อ ack หรือ return error เพื่อ dead-letter ส่วน body ที่ไม่มีวัน parse ได้คือเคส poison message ที่การ dead-letter ครั้งเดียวเหมาะพอดี
การแยกการส่งออกมาเป็น process ของตัวเองคือผลตอบแทนทั้งหมดของการ enqueue แทนการส่ง inline งานส่งที่ช้าและพังได้จะ crash, retry และ scale ผ่าน competing consumers ได้ตามลำพัง — พิสูจน์แล้วด้วย worker สองตัวแบ่ง burst หก job — โดยไม่เคยไปหยุด Kafka consumer ที่ต้องเร็วและคอยเติม queue อยู่
การที่ทั้งสอง hop เป็น at-least-once แปลว่า Deliver ระดับ production จะต้อง dedupe บน order_id + outcome ก่อนส่งจริง แบบเดียวกับ persisted idempotency ที่การ charge card จริงต้องมี เท่านี้เซอร์วิส Notification ก็สมบูรณ์ พร้อมกับสถาปัตยกรรมสอง broker แบบ end to end: REST request เดียวไหลผ่านห้าเซอร์วิสและผ่านทั้ง Kafka และ RabbitMQ โดยไม่มี synchronous coupling ที่ไหนเลย บทถัดไป Resilience → ถอยจาก feature มาที่การ hardening — timeout, retry พร้อม backoff, และ circuit breaker ข้ามทุก hop ที่ระบบนี้มีตอนนี้