1. Message Queue Overview
| Feature | RabbitMQ | Kafka | NATS |
|---|---|---|---|
| Model | Message Broker | Event Log | Pub/Sub |
| Ordering | Per queue | Per partition | No guarantee |
| Persistence | Optional | Always (retention) | Optional (JetStream) |
| Throughput | ~50K msg/s | ~1M msg/s | ~10M msg/s |
| Use case | Task queues, RPC | Event streaming | Microservices, IoT |
| Complexity | Medium | High | Low |
2. RabbitMQ với Go
2.1. Setup
go get github.com/rabbitmq/amqp091-go
2.2. Connection Manager
package rabbitmq
import (
"context"
"fmt"
"log"
"sync"
"time"
amqp "github.com/rabbitmq/amqp091-go"
)
type Connection struct {
conn *amqp.Connection
channel *amqp.Channel
mu sync.RWMutex
url string
done chan struct{}
}
func NewConnection(url string) (*Connection, error) {
c := &Connection{url: url, done: make(chan struct{})}
if err := c.connect(); err != nil {
return nil, err
}
go c.reconnectLoop()
return c, nil
}
func (c *Connection) connect() error {
conn, err := amqp.Dial(c.url)
if err != nil {
return fmt.Errorf("dial: %w", err)
}
ch, err := conn.Channel()
if err != nil {
conn.Close()
return fmt.Errorf("channel: %w", err)
}
c.mu.Lock()
c.conn = conn
c.channel = ch
c.mu.Unlock()
return nil
}
func (c *Connection) reconnectLoop() {
for {
select {
case <-c.done:
return
case err := <-c.conn.NotifyClose(make(chan *amqp.Error)):
if err != nil {
log.Printf("RabbitMQ connection lost: %v", err)
}
for i := 0; i < 30; i++ {
time.Sleep(time.Duration(i+1) * time.Second)
if err := c.connect(); err == nil {
log.Println("RabbitMQ reconnected")
break
}
}
}
}
}
func (c *Connection) Close() {
close(c.done)
c.channel.Close()
c.conn.Close()
}
2.3. Publisher
type Publisher struct {
conn *Connection
}
func NewPublisher(conn *Connection) *Publisher {
return &Publisher{conn: conn}
}
func (p *Publisher) Publish(ctx context.Context, exchange, routingKey string, body []byte) error {
p.conn.mu.RLock()
ch := p.conn.channel
p.conn.mu.RUnlock()
return ch.PublishWithContext(ctx,
exchange,
routingKey,
false, // mandatory
false, // immediate
amqp.Publishing{
ContentType: "application/json",
Body: body,
DeliveryMode: amqp.Persistent,
Timestamp: time.Now(),
MessageId: uuid.New().String(),
},
)
}
// Setup exchange & queue
func SetupExchange(ch *amqp.Channel) error {
// Topic exchange cho event routing
if err := ch.ExchangeDeclare(
"events", // name
"topic", // type: direct, fanout, topic, headers
true, // durable
false, // auto-deleted
false, // internal
false, // no-wait
nil, // arguments
); err != nil {
return err
}
// Dead letter exchange
if err := ch.ExchangeDeclare(
"events.dlx",
"topic",
true, false, false, false, nil,
); err != nil {
return err
}
return nil
}
func SetupQueue(ch *amqp.Channel, queueName, exchange, routingKey string) error {
// Main queue with DLX
q, err := ch.QueueDeclare(
queueName,
true, // durable
false, // auto-delete
false, // exclusive
false, // no-wait
amqp.Table{
"x-dead-letter-exchange": "events.dlx",
"x-dead-letter-routing-key": queueName + ".dead",
"x-message-ttl": int32(86400000), // 24h
},
)
if err != nil {
return err
}
return ch.QueueBind(q.Name, routingKey, exchange, false, nil)
}
2.4. Consumer
type Consumer struct {
conn *Connection
handlers map[string]MessageHandler
}
type MessageHandler func(ctx context.Context, body []byte) error
func NewConsumer(conn *Connection) *Consumer {
return &Consumer{
conn: conn,
handlers: make(map[string]MessageHandler),
}
}
func (c *Consumer) Register(eventType string, handler MessageHandler) {
c.handlers[eventType] = handler
}
func (c *Consumer) Start(ctx context.Context, queueName string, concurrency int) error {
c.conn.mu.RLock()
ch := c.conn.channel
c.conn.mu.RUnlock()
// Prefetch: limit unacked messages
if err := ch.Qos(concurrency, 0, false); err != nil {
return err
}
msgs, err := ch.Consume(
queueName,
"", // consumer tag
false, // auto-ack (manual!)
false, // exclusive
false, // no-local
false, // no-wait
nil,
)
if err != nil {
return err
}
// Worker pool
sem := make(chan struct{}, concurrency)
for {
select {
case <-ctx.Done():
return ctx.Err()
case msg, ok := <-msgs:
if !ok {
return fmt.Errorf("channel closed")
}
sem <- struct{}{}
go func(d amqp.Delivery) {
defer func() { <-sem }()
if err := c.processMessage(ctx, d); err != nil {
log.Printf("Process error: %v", err)
d.Nack(false, false) // Reject → DLQ
} else {
d.Ack(false)
}
}(msg)
}
}
}
func (c *Consumer) processMessage(ctx context.Context, d amqp.Delivery) error {
eventType := d.RoutingKey
handler, ok := c.handlers[eventType]
if !ok {
log.Printf("No handler for event: %s", eventType)
return nil // Ack unknown events
}
return handler(ctx, d.Body)
}
3. Kafka với Go
go get github.com/segmentio/kafka-go
package kafka
import (
"context"
"encoding/json"
"log"
"time"
"github.com/segmentio/kafka-go"
)
// Producer
type KafkaProducer struct {
writer *kafka.Writer
}
func NewProducer(brokers []string, topic string) *KafkaProducer {
return &KafkaProducer{
writer: &kafka.Writer{
Addr: kafka.TCP(brokers...),
Topic: topic,
Balancer: &kafka.LeastBytes{},
BatchSize: 100,
BatchTimeout: 10 * time.Millisecond,
RequiredAcks: kafka.RequireOne,
},
}
}
func (p *KafkaProducer) Publish(ctx context.Context, key string, value interface{}) error {
data, err := json.Marshal(value)
if err != nil {
return err
}
return p.writer.WriteMessages(ctx, kafka.Message{
Key: []byte(key),
Value: data,
Time: time.Now(),
})
}
func (p *KafkaProducer) Close() error {
return p.writer.Close()
}
// Consumer Group
type KafkaConsumer struct {
reader *kafka.Reader
}
func NewConsumer(brokers []string, topic, groupID string) *KafkaConsumer {
return &KafkaConsumer{
reader: kafka.NewReader(kafka.ReaderConfig{
Brokers: brokers,
Topic: topic,
GroupID: groupID,
MinBytes: 1,
MaxBytes: 10e6, // 10MB
CommitInterval: time.Second,
StartOffset: kafka.LastOffset,
}),
}
}
func (c *KafkaConsumer) Consume(ctx context.Context, handler func(kafka.Message) error) error {
for {
msg, err := c.reader.ReadMessage(ctx)
if err != nil {
if ctx.Err() != nil {
return nil
}
log.Printf("Read error: %v", err)
continue
}
if err := handler(msg); err != nil {
log.Printf("Handle error (offset %d): %v", msg.Offset, err)
}
}
}
4. Event-Driven Patterns
4.1. Domain Events
// Domain event
type Event struct {
ID string `json:"id"`
Type string `json:"type"`
Source string `json:"source"`
Data any `json:"data"`
Timestamp time.Time `json:"timestamp"`
}
// Event types
const (
UserCreated = "user.created"
UserUpdated = "user.updated"
OrderCreated = "order.created"
OrderPaid = "order.paid"
)
// Event bus interface
type EventBus interface {
Publish(ctx context.Context, event Event) error
Subscribe(eventType string, handler func(Event) error)
}
4.2. Outbox Pattern (Transactional Outbox)
// Outbox table
type Outbox struct {
ID uint `gorm:"primaryKey"`
EventType string `gorm:"index"`
Payload string `gorm:"type:jsonb"`
Published bool `gorm:"default:false;index"`
CreatedAt time.Time
}
// Save event in same transaction as business data
func (s *OrderService) CreateOrder(ctx context.Context, input CreateOrderInput) error {
return s.db.Transaction(func(tx *gorm.DB) error {
// 1. Save order
order := &Order{UserID: input.UserID, Total: input.Total}
if err := tx.Create(order).Error; err != nil {
return err
}
// 2. Save event to outbox (same transaction!)
event := Event{
ID: uuid.New().String(),
Type: OrderCreated,
Data: order,
Timestamp: time.Now(),
}
payload, _ := json.Marshal(event)
outbox := &Outbox{
EventType: OrderCreated,
Payload: string(payload),
}
return tx.Create(outbox).Error
})
}
// Outbox publisher (separate goroutine)
func (s *OutboxPublisher) PollAndPublish(ctx context.Context) {
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
var events []Outbox
s.db.Where("published = ?", false).
Order("id ASC").
Limit(100).
Find(&events)
for _, e := range events {
if err := s.publisher.Publish(ctx, "events", e.EventType, []byte(e.Payload)); err != nil {
log.Printf("Publish error: %v", err)
continue
}
s.db.Model(&e).Update("published", true)
}
}
}
}
5. Retry Strategies
type RetryConfig struct {
MaxRetries int
InitDelay time.Duration
MaxDelay time.Duration
Multiplier float64
}
func WithRetry(ctx context.Context, cfg RetryConfig, fn func() error) error {
delay := cfg.InitDelay
for attempt := 0; attempt <= cfg.MaxRetries; attempt++ {
err := fn()
if err == nil {
return nil
}
if attempt == cfg.MaxRetries {
return fmt.Errorf("max retries exceeded: %w", err)
}
log.Printf("Attempt %d failed: %v. Retrying in %v", attempt+1, err, delay)
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(delay):
}
// Exponential backoff
delay = time.Duration(float64(delay) * cfg.Multiplier)
if delay > cfg.MaxDelay {
delay = cfg.MaxDelay
}
}
return nil
}
Bài tiếp theo: Caching, File Upload & Performance — Redis caching, file upload, và optimization.