前言
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)