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

Producer & Consumer

pkg/events/event.go — struct Event เดียวที่เป็น JSON envelope ของทุกข้อความที่ระบบนี้ publish ไป Kafka ไม่ว่าเซอร์วิสไหนจะเป็นคน produce หรือ consume ก็ตาม และ pkg/kafka/kafka.goPublisher กับ Consumer ทั้งคู่เป็น wrapper บาง ๆ รอบ segmentio/kafka-go เพื่อให้ทุกเซอร์วิส produce และ consume envelope นั้นด้วยวิธีเดียวกัน แทนที่แต่ละตัวจะไปตั้งค่า kafka.Writer/kafka.Reader ของตัวเองใหม่ตั้งแต่ต้น ทั้งสอง package เป็นโค้ดที่ใช้ร่วมกัน: The Transactional Outbox → ระบุชื่อ relay ที่จะใช้ Publisher ไว้แล้ว (สร้างในบทถัดไป) และ Payment, Notification, และ Order Saga — Module 8, 9, และ 10 — จะใช้ Consumer เพื่อตอบสนองต่อ event เดียวกันเหล่านี้ทั้งหมด

เพื่อพิสูจน์ว่า wrapper นี้ใช้งานได้จริงกับ broker จริงก่อนจะมีใครมาพึ่งพา บทนี้ยังเขียนโปรแกรมเล็ก ๆ แบบใช้แล้วทิ้งสองตัว — cmd/kafkademo/produce กับ cmd/kafkademo/consume — ที่ publish event หนึ่งตัวแล้ว consume กลับมา กับ container Kafka ที่ Infra & Compose → สร้างไว้แล้วใน Module 1

ทุกเซอร์วิสในระบบนี้สุดท้ายต้องทั้ง publish หรือ consume ข้อความ Kafka — Order (โมดูลนี้), Payment, Notification, และ Order Saga ทำทั้งหมด ถ้าไม่มี envelope ร่วมกัน แต่ละเซอร์วิสจะคิดรูปแบบ JSON ของ “สิ่งที่เกิดขึ้น” ขึ้นมาเอง แล้ว consumer ที่ทีมหนึ่งเขียนก็ต้องรู้ quirk เฉพาะของรูปแบบที่ producer ของอีกทีมเลือกใช้ events.Event แก้ปัญหานี้ด้วย struct เดียว รูปแบบ JSON เดียว ใช้เหมือนกันทุกที่ ID ระบุ event นั้นแบบไม่ซ้ำกัน (และใน Outbox & Relay → จะกลายเป็น key ที่ consumer ที่ idempotent ใช้ dedupe), Type คือชื่อ event (order.created, payment.succeeded และอื่น ๆ), AggregateID คือ id ของ entity ที่ event นี้พูดถึง, Payload คือข้อมูลเฉพาะของ event นั้นในรูป JSON ดิบ และ OccurredAt คือเวลาที่ event เกิดขึ้น

เมธอด Run ของ Consumer commit offset ด้วยมือ — CommitMessages จะรันก็ต่อเมื่อ handle คืนค่าสำเร็จเท่านั้น ไม่เคยรันก่อนหน้านั้นเลย การจัดลำดับแบบนี้เพียงอย่างเดียวคือสิ่งที่ทำให้ consumer ของ Kafka ในระบบนี้เป็น at-least-once ไม่ใช่ at-most-once: ถ้า consumer crash ระหว่างที่ handle ทำเสร็จแล้วกับตอนที่ commit ไปถึง Kafka ไม่มีบันทึกเลยว่าข้อความนี้ถูกประมวลผลแล้ว แล้ว Kafka จะ redeliver อีกครั้งหลัง restart นี่คือการตัดสินใจโดยตั้งใจ ไม่ใช่ความผิดพลาด ทางเลือกอีกแบบคือ commit ก่อนแล้วค่อยประมวลผล ที่เป็น at-most-once และการ crash ระหว่างสองขั้นตอนนั้นจะทำให้ event หายไปตลอดกาลอย่างเงียบ ๆ การทำ order.created หรือ payment.succeeded หายไปแย่กว่าการประมวลผลซ้ำเป็นบางครั้งมาก นี่คือเหตุผลที่ consumer ทุกตัวที่สร้างบน Consumer.Run ต้องเป็น idempotent — ข้อกำหนดที่ Architecture → ระบุไว้ตั้งแต่โมดูลแรก

pkg/kafka wrapper ที่ใช้ร่วมกัน เทียบกับ แต่ละเซอร์วิสตั้งค่า kafka.Writer/kafka.Reader เอง

  • Pros: ทุก producer ได้ Balancer, RequiredAcks, และพฤติกรรมสร้าง topic แบบเดียวกันโดยไม่ต้อง copy-paste config ไปห้าเซอร์วิส; ทุก consumer ได้ loop fetch → handle → commit แบบเดียวกัน ดังนั้นคำถาม “จะรับประกัน at-least-once ตรงนี้ยังไง” ถูกตอบครั้งเดียวในไฟล์เดียว แทนที่จะตอบทีละเซอร์วิส (พร้อมความเสี่ยงที่บางเซอร์วิสจะทำผิดแบบแนบเนียน)
  • Cons: abstraction ที่ใช้ร่วมกันตัวเดียวตอนนี้ต้อง general พอสำหรับทุก consumer — Consumer.Run รับ callback handle แทนที่จะ expose API เต็มของ Reader จาก kafka-go ดังนั้นเซอร์วิสที่ต้องการฟีเจอร์ที่ wrapper ไม่ได้ expose จริง ๆ (เช่น callback สำหรับ partition assignment) ต้องขยาย wrapper เอง หรือใช้ kafka-go ตรง ๆ เฉพาะกรณีนั้น

Manual commit (CommitMessages หลังจาก handle สำเร็จ) เทียบกับ auto-commit ในตัวของ kafka-go (CommitInterval บน ReaderConfig)

  • Pros: commit จะเกิดขึ้นก็ต่อเมื่อข้อความถูกจัดการเสร็จสมบูรณ์แล้วเท่านั้น ดังนั้น crash กลางทางระหว่างประมวลผลจะจบด้วยการ redeliver เสมอ ไม่มีทางหายไปเงียบ ๆ — นี่คือการรับประกันที่ข้อกำหนด idempotency ของระบบนี้ถูกสร้างขึ้นมาบนฐานนั้น
  • Cons: consumer ทุกตัวตอนนี้ต้องเป็น idempotent จริง ๆ เพราะข้อความเดียวกันมาถึงสองครั้งได้จริง; auto-commit เข้าใจง่ายกว่าสำหรับ handler ที่ไม่สนใจความต่างนี้ (แลกกับช่วงเวลาสั้น ๆ แบบ at-most-once ที่ offset commit วิ่งแซงหน้าการประมวลผลไปก่อน)
// Package events defines the JSON envelope every Kafka message in this
// system carries. Every producer marshals one, every consumer unmarshals
// one — no service ever hand-rolls its own message shape.
package events
import (
"encoding/json"
"time"
)
// Event is the JSON envelope of every message published to any topic in
// this system.
type Event struct {
ID string `json:"id"`
Type string `json:"type"`
AggregateID string `json:"aggregate_id"`
Payload json.RawMessage `json:"payload"`
OccurredAt time.Time `json:"occurred_at"`
}

เซฟไฟล์นี้เป็น pkg/events/event.go Payload เป็น json.RawMessage ไม่ใช่ struct ที่เจาะจง — events.Event ไม่รู้และไม่สนใจว่า payload ของ event ตัวไหนหน้าตาเป็นอย่างไร เพียงแค่พา byte ที่ marshal ไว้แล้วส่งผ่านไป จากนั้นให้แต่ละ consumer unmarshal Payload เป็น struct ที่ event type นั้นต้องการเอง

// Package kafka wraps segmentio/kafka-go's Writer and Reader behind a
// small, project-specific Publisher/Consumer pair, so every service
// produces and consumes the events.Event envelope the same way.
package kafka
import (
"context"
"encoding/json"
"fmt"
"log"
"github.com/avetavos/shopmicro/pkg/events"
kafkago "github.com/segmentio/kafka-go"
)
// Publisher writes messages to Kafka.
type Publisher struct {
w *kafkago.Writer
}
// NewPublisher returns a Publisher connected to brokers. Balancer: &Hash{}
// routes messages sharing the same key to the same partition;
// AllowAutoTopicCreation lets a topic come into existence on first publish;
// RequiredAcks: RequireAll waits for every in-sync replica to acknowledge
// the write before WriteMessages returns.
func NewPublisher(brokers []string) *Publisher {
return &Publisher{
w: &kafkago.Writer{
Addr: kafkago.TCP(brokers...),
Balancer: &kafkago.Hash{},
AllowAutoTopicCreation: true,
RequiredAcks: kafkago.RequireAll,
},
}
}
// Publish writes value to topic, keyed by key. Messages with the same key
// always land on the same partition, which is why callers key by an
// aggregate's id (an order's id, say) — every event about that aggregate
// stays in the one place Kafka guarantees ordering: within a partition.
func (p *Publisher) Publish(ctx context.Context, topic, key string, value []byte) error {
if err := p.w.WriteMessages(ctx, kafkago.Message{
Topic: topic,
Key: []byte(key),
Value: value,
}); err != nil {
return fmt.Errorf("kafka: publish to %s: %w", topic, err)
}
return nil
}
// Close flushes any pending writes and closes the underlying connection.
func (p *Publisher) Close() error {
return p.w.Close()
}
// Consumer reads events.Event envelopes from a single topic as part of a
// consumer group.
type Consumer struct {
r *kafkago.Reader
}
// NewConsumer returns a Consumer that reads topic as part of groupID.
// Give every service its own groupID (its service name is enough) — Kafka
// delivers every message on topic to every distinct group at least once,
// so each service sees the full stream regardless of what any other
// service's consumers are doing. Multiple processes sharing the same
// groupID split the topic's partitions between them instead.
func NewConsumer(brokers []string, groupID, topic string) *Consumer {
return &Consumer{
r: kafkago.NewReader(kafkago.ReaderConfig{
Brokers: brokers,
GroupID: groupID,
Topic: topic,
MinBytes: 10e3, // 10KB
MaxBytes: 10e6, // 10MB
}),
}
}
// Run fetches messages from topic in a loop, unmarshals each into an
// events.Event, and passes it to handle. A message's offset is committed
// with CommitMessages only after handle returns nil — this is manual-commit,
// at-least-once delivery: if the process crashes after handle succeeds but
// before the commit lands, Kafka redelivers the same message on restart, so
// handle must be idempotent. If handle (or unmarshalling) fails, Run logs
// the error and moves on without committing, so the same message is
// redelivered on the next FetchMessage rather than silently dropped. Run
// blocks until ctx is cancelled, at which point FetchMessage returns ctx's
// error.
func (c *Consumer) Run(ctx context.Context, handle func(context.Context, events.Event) error) error {
for {
msg, err := c.r.FetchMessage(ctx)
if err != nil {
return fmt.Errorf("kafka: fetch message: %w", err)
}
var ev events.Event
if err := json.Unmarshal(msg.Value, &ev); err != nil {
log.Printf("kafka: unmarshal event at %s/%d/%d: %v", msg.Topic, msg.Partition, msg.Offset, err)
continue
}
if err := handle(ctx, ev); err != nil {
log.Printf("kafka: handle event %s (%s): %v — will redeliver", ev.ID, ev.Type, err)
continue
}
if err := c.r.CommitMessages(ctx, msg); err != nil {
return fmt.Errorf("kafka: commit message: %w", err)
}
}
}
// Close closes the underlying reader.
func (c *Consumer) Close() error {
return c.r.Close()
}

เซฟไฟล์นี้เป็น pkg/kafka/kafka.go แล้วดึง client library เข้ามา:

Terminal window
go get github.com/segmentio/kafka-go

โปรแกรมใช้แล้วทิ้งสองตัว ไม่ได้เป็นส่วนหนึ่งของเซอร์วิสไหนเลย มีไว้เพื่อพิสูจน์ว่า Publisher กับ Consumer ทำงานได้จริงกับ broker จริงเท่านั้น

// Command produce publishes a single demo events.Event to the "orders"
// topic — nothing wired into any service, just proof that Publisher works
// end to end against a running Kafka broker.
package main
import (
"context"
"encoding/json"
"log"
"strings"
"time"
"github.com/avetavos/shopmicro/pkg/config"
"github.com/avetavos/shopmicro/pkg/events"
"github.com/avetavos/shopmicro/pkg/kafka"
)
func main() {
brokers := strings.Split(config.Get("KAFKA_BROKERS", "localhost:9092"), ",")
pub := kafka.NewPublisher(brokers)
defer pub.Close()
ev := events.Event{
ID: "demo-1",
Type: "order.created",
AggregateID: "order-abc",
Payload: json.RawMessage(`{"order_id":"order-abc","total_cents":2598}`),
OccurredAt: time.Now().UTC(),
}
value, err := json.Marshal(ev)
if err != nil {
log.Fatalf("marshal event: %v", err)
}
if err := pub.Publish(context.Background(), "orders", ev.AggregateID, value); err != nil {
log.Fatalf("publish: %v", err)
}
log.Println("published:", ev.ID)
}

เซฟไฟล์นี้เป็น cmd/kafkademo/produce/main.go

// Command consume reads events.Event messages from the "orders" topic and
// prints each one — nothing wired into any service, just proof that
// Consumer.Run's fetch/handle/commit loop works end to end.
package main
import (
"context"
"log"
"strings"
"github.com/avetavos/shopmicro/pkg/config"
"github.com/avetavos/shopmicro/pkg/events"
"github.com/avetavos/shopmicro/pkg/kafka"
)
func main() {
brokers := strings.Split(config.Get("KAFKA_BROKERS", "localhost:9092"), ",")
consumer := kafka.NewConsumer(brokers, "kafkademo", "orders")
defer consumer.Close()
err := consumer.Run(context.Background(), func(ctx context.Context, ev events.Event) error {
log.Printf("consumed: id=%s type=%s aggregate_id=%s payload=%s", ev.ID, ev.Type, ev.AggregateID, string(ev.Payload))
return nil
})
if err != nil {
log.Fatalf("consumer: %v", err)
}
}

เซฟไฟล์นี้เป็น cmd/kafkademo/consume/main.go ทั้งสองไฟล์ default เป็น KAFKA_BROKERS=localhost:9092 ตัวแปรและค่าเริ่มต้นเดียวกับที่ .env.example ของ Repo Layout → ระบุไว้แล้ว

ตรวจสอบว่า container Kafka ของ Module 1 รันอยู่:

Terminal window
cd deploy/compose && docker compose up -d kafka

เทอร์มินัลแรก ให้รัน consumer ก่อน เพื่อให้ subscribe รออยู่แล้วตอนที่ message ถูก publish:

Terminal window
go run ./cmd/kafkademo/consume

เทอร์มินัลที่สอง publish demo event:

Terminal window
go run ./cmd/kafkademo/produce
published: demo-1

กลับไปที่เทอร์มินัลของ consumer:

consumed: id=demo-1 type=order.created aggregate_id=order-abc payload={"order_id":"order-abc","total_cents":2598}

AllowAutoTopicCreation หมายความว่า topic orders ไม่ต้องมีอยู่มาก่อน เพราะการ publish ครั้งแรกจะสร้างขึ้นมาเอง หยุด consumer ด้วย Ctrl-C

จากนั้นยืนยันว่า module ยัง build ผ่าน:

Terminal window
go build ./...

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

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

  • ทำไม Event.Payload ถึงใช้ json.RawMessage แทนที่จะเป็น struct ที่เจาะจง?
  • ถ้า Consumer.Run เรียก CommitMessages ก่อน เรียก handle แทนที่จะเรียกทีหลัง จะเปลี่ยนอะไรบ้าง?
  • ถ้า handle คืนค่า error ทำไม Run ถึงข้ามการ commit แล้วไปต่อ แทนที่จะคืน error ทันทีแล้วหยุด consumer ทั้งตัว?

pkg/events.Event คือ JSON envelope เดียวที่ทุก message บน Kafka ในระบบนี้ใช้ — ID, Type, AggregateID, Payload, OccurredAt — ดังนั้นไม่มีสองเซอร์วิสไหนคิดรูปแบบของตัวเองที่เข้ากันไม่ได้สำหรับ “สิ่งที่เกิดขึ้น” pkg/kafka.Publisher ห่อ kafka.Writer ที่ตั้งค่าด้วย &kafka.Hash{} (partition ตาม key), AllowAutoTopicCreation, และ RequiredAcks: RequireAll pkg/kafka.Consumer ห่อ kafka.Reader แล้ว expose Run ซึ่ง loop fetch → handle → commit จะเรียก CommitMessages ก็ต่อเมื่อ handle สำเร็จเท่านั้น — manual-commit, delivery แบบ at-least-once นั่นคือเหตุผลว่าทำไม consumer ทุกตัวที่สร้างบนนี้ต้องเป็น idempotent cmd/kafkademo/produce กับ cmd/kafkademo/consume พิสูจน์แล้วว่าทั้งสองฝั่งทำงานได้กับ Kafka broker จริงจาก Module 1 publish แล้ว consume event หนึ่งตัวแบบ end-to-end ต่อไป Topics, Partitions & Consumer Groups → จะอธิบายแนวคิดที่บทนี้ใช้ไปโดยไม่ได้เอ่ยชื่อ — partition, key, consumer group, และความ replay ได้ของ log — ก่อนที่ Outbox & Relay → จะเอา Publisher ไปใช้งานจริง ปิด outbox pattern ที่ The Transactional Outbox → เริ่มไว้