Go + RabbitMQ 构建高并发任务队列实战

Go + RabbitMQ 构建高并发任务队列实战

admin
2026-07-24 / 0 评论 / 0 阅读

前言

Go 语言的 goroutine 并发模型配合 RabbitMQ 的可靠消息投递,是构建高并发任务队列的经典方案。本文将实现一个完整的 Go + RabbitMQ 任务队列系统,包含生产者、消费者、重试机制和优雅退出。

一、项目结构

go-rabbitmq-queue/
├── go.mod
├── config/
│   └── config.go
├── producer/
│   └── main.go
├── consumer/
│   └── main.go
└── rabbitmq/
    └── connection.go

二、安装依赖

go mod init github.com/yourname/go-rabbitmq-queue
go get github.com/rabbitmq/amqp091-go

三、RabbitMQ 连接封装

package rabbitmq

import (
    "log"
    "time"

    amqp "github.com/rabbitmq/amqp091-go"
)

type Config struct {
    URL           string
    Exchange      string
    Queue         string
    RoutingKey    string
    PrefetchCount int
}

type RabbitMQ struct {
    conn    *amqp.Connection
    channel *amqp.Channel
    config  Config
}

func New(cfg Config) (*RabbitMQ, error) {
    var conn *amqp.Connection
    var err error

    for i := 0; i < 5; i++ {
        conn, err = amqp.Dial(cfg.URL)
        if err == nil {
            break
        }
        log.Printf("连接失败(%d/5): %v", i+1, err)
        time.Sleep(3 * time.Second)
    }
    if err != nil {
        return nil, err
    }

    ch, err := conn.Channel()
    if err != nil {
        return nil, err
    }

    err = ch.ExchangeDeclare(
        cfg.Exchange, "direct",
        true, false, false, false, nil,
    )
    if err != nil {
        return nil, err
    }

    args := amqp.Table{
        "x-message-ttl":          int32(60000),
        "x-dead-letter-exchange": cfg.Exchange + ".dlx",
    }

    _, err = ch.QueueDeclare(
        cfg.Queue, true, false, false, false, args,
    )
    if err != nil {
        return nil, err
    }

    err = ch.QueueBind(
        cfg.Queue, cfg.RoutingKey, cfg.Exchange, false, nil,
    )
    if err != nil {
        return nil, err
    }

    ch.Qos(cfg.PrefetchCount, 0, false)

    return &RabbitMQ{conn: conn, channel: ch, config: cfg}, nil
}

func (r *RabbitMQ) Channel() *amqp.Channel {
    return r.channel
}

func (r *RabbitMQ) Close() {
    r.channel.Close()
    r.conn.Close()
}

四、生产者

package main

import (
    "encoding/json"
    "fmt"
    "log"
    "time"

    "github.com/yourname/go-rabbitmq-queue/rabbitmq"
    amqp "github.com/rabbitmq/amqp091-go"
)

type Task struct {
    ID      string      `json:"id"`
    Type    string      `json:"type"`
    Payload interface{} `json:"payload"`
}

func main() {
    cfg := rabbitmq.Config{
        URL:           "amqp://admin:admin123@localhost:5672/",
        Exchange:      "task.exchange",
        Queue:         "task.queue",
        RoutingKey:    "task.process",
        PrefetchCount: 10,
    }

    mq, err := rabbitmq.New(cfg)
    if err != nil {
        log.Fatal(err)
    }
    defer mq.Close()

    for i := 0; i < 100; i++ {
        task := Task{
            ID:   fmt.Sprintf("task-%d", i),
            Type: "email",
            Payload: map[string]string{
                "to":      fmt.Sprintf("user%d@example.com", i),
                "subject": "通知邮件",
            },
        }

        body, _ := json.Marshal(task)

        err = mq.Channel().Publish(
            cfg.Exchange, cfg.RoutingKey,
            false, false,
            amqp.Publishing{
                DeliveryMode: amqp.Persistent,
                ContentType:  "application/json",
                Body:         body,
                Timestamp:    time.Now(),
            },
        )
        if err != nil {
            log.Printf("发送失败: %v", err)
            continue
        }

        log.Printf("发送任务: %s", task.ID)
    }

    log.Println("所有任务发送完成")
}

五、消费者(多 Worker 并发)

package main

import (
    "encoding/json"
    "fmt"
    "log"
    "os"
    "os/signal"
    "sync"
    "syscall"
    "time"

    "github.com/yourname/go-rabbitmq-queue/rabbitmq"
    amqp "github.com/rabbitmq/amqp091-go"
)

type Task struct {
    ID      string      `json:"id"`
    Type    string      `json:"type"`
    Payload interface{} `json:"payload"`
}

func processTask(task Task) error {
    log.Printf("处理任务: %s, 类型: %s", task.ID, task.Type)
    time.Sleep(500 * time.Millisecond)

    if time.Now().Unix()%10 == 0 {
        return fmt.Errorf("模拟处理失败")
    }

    log.Printf("任务完成: %s", task.ID)
    return nil
}

func startWorker(id int, mq *rabbitmq.RabbitMQ, wg *sync.WaitGroup) {
    defer wg.Done()

    msgs, err := mq.Channel().Consume(
        "task.queue",
        fmt.Sprintf("worker-%d", id),
        false, false, false, false, nil,
    )
    if err != nil {
        log.Printf("Worker %d 启动失败: %v", id, err)
        return
    }

    log.Printf("Worker %d 启动", id)

    for msg := range msgs {
        var task Task
        if err := json.Unmarshal(msg.Body, &task); err != nil {
            log.Printf("Worker %d 解析失败: %v", id, err)
            msg.Nack(false, false)
            continue
        }

        if err := processTask(task); err != nil {
            log.Printf("Worker %d 处理失败: %s -> %v", id, task.ID, err)
            retryCount := getRetryCount(msg)
            if retryCount < 3 {
                msg.Nack(false, true)
            } else {
                log.Printf("Worker %d 任务 %s 重试超限,进入死信", id, task.ID)
                msg.Nack(false, false)
            }
            continue
        }

        msg.Ack(false)
    }

    log.Printf("Worker %d 退出", id)
}

func getRetryCount(msg amqp.Delivery) int {
    if deaths, ok := msg.Headers["x-death"].([]interface{}); ok && len(deaths) > 0 {
        if death, ok := deaths[0].(amqp.Table); ok {
            if count, ok := death["count"].(int64); ok {
                return int(count)
            }
        }
    }
    return 0
}

func main() {
    cfg := rabbitmq.Config{
        URL:           "amqp://admin:admin123@localhost:5672/",
        Exchange:      "task.exchange",
        Queue:         "task.queue",
        RoutingKey:    "task.process",
        PrefetchCount: 5,
    }

    mq, err := rabbitmq.New(cfg)
    if err != nil {
        log.Fatal(err)
    }
    defer mq.Close()

    var wg sync.WaitGroup
    workerCount := 5
    for i := 1; i <= workerCount; i++ {
        wg.Add(1)
        go startWorker(i, mq, &wg)
    }

    sigs := make(chan os.Signal, 1)
    signal.Notify(sigs, syscall.SIGINT, syscall.SIGTERM)
    <-sigs

    log.Println("收到退出信号,等待 worker 完成...")
    mq.Close()
    wg.Wait()
    log.Println("所有 worker 已退出")
}

六、死信队列消费者

func startDeadLetterConsumer(mq *rabbitmq.RabbitMQ) {
    ch := mq.Channel()
    ch.ExchangeDeclare("task.exchange.dlx", "fanout", true, false, false, false, nil)
    _, _ = ch.QueueDeclare("task.dlq", true, false, false, false, nil)
    _ = ch.QueueBind("task.dlq", "", "task.exchange.dlx", false, nil)

    msgs, _ := ch.Consume("task.dlq", "dlq-consumer", false, false, false, false, nil)

    go func() {
        for msg := range msgs {
            log.Printf("死信消息: %s", string(msg.Body))
            msg.Ack(false)
        }
    }()
}

七、架构总结

Producer → [task.exchange] → [task.queue] → 5个 Worker 并发消费
                                    ↓ (失败/Nack)
                              [task.exchange.dlx] → [task.dlq] → 死信消费者记录

总结

Go + RabbitMQ 的组合充分发挥了各自优势:Go 的 goroutine 让消费者可以轻松开几十个并发 worker,RabbitMQ 的 ACK 机制保证消息不丢失。生产环境注意:设置合理的 prefetch 避免消息堆积在单个 worker、实现死信队列处理失败消息、添加优雅退出逻辑避免消息中断。

0

评论 (0)

取消
0:00