Chuyển đến nội dung chính

レッスン 15: メッセージ キューとイベント駆動型

Go を使用した RabbitMQ、交換、キュー。 Kafka のプロデューサー/コンシューマー。イベント駆動型のアーキテクチャ パターン。サーガパターン、アウトボックスパターン。デッドレターキューと再試行戦略。

💻 プログラミング — レッスン 15 レッスン 15: メッセージ キューとイベント駆動型

Golang: 基本から高度まで

パート 4: 高度な機能

xdev.asia

1. メッセージキューの概要

特長ラビットMQカフカNATS
モデルメッセージブローカーイベントログパブ/サブ
注文キューごとパーティションごと保証なし
持続性オプション常に(保持)オプション(ジェットストリーム)
スループット~50,000 メッセージ/秒~100万メッセージ/秒~1,000万メッセージ/秒
ユースケースタスクキュー、RPCイベントストリーミングマイクロサービス、IoT
複雑さ中高低い

2. Go を使用した RabbitMQ

2.1.セットアップ

go get github.com/rabbitmq/amqp091-go

2.2.接続マネージャー

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.出版社

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.消費者

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. カフカと 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. イベント駆動型パターン

4.1.ドメインイベント

// 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 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. 再試行戦略

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
}

次の記事: キャッシュ、ファイルのアップロード、パフォーマンス — Redis キャッシュ、ファイルのアップロード、最適化。