返回首页

Golang项目实践:4.5 消息队列与异步处理

在微服务、订单系统、支付系统、通知系统等场景中,**消息队列(Message Queue)**是实现系统解耦、削峰填谷、异步处理和事件驱动的重要基础设施。对于 Golang 开发者来说,掌握 Kafka、RabbitMQ、NATS 等常见消息中间件的使用方式,以及消息可靠性、幂等、延迟任务等工程实践,是构建高可用后端系统的关键能力。

本文将围绕以下几个方面展开:

  • Kafka 生产者与消费者实现
  • RabbitMQ / NATS 集成
  • 事件驱动架构
  • 消息可靠性保证与幂等设计
  • 延迟队列与定时任务

文章中的示例均采用 Go 编写,并尽量提供可直接运行的代码。

一、为什么需要消息队列

在同步调用模型里,一个请求往往要串行完成多个动作,例如:

  1. 创建订单
  2. 扣减库存
  3. 发送短信
  4. 写审计日志
  5. 触发推荐系统更新

如果这些动作全部放在一个同步链路中,会带来几个明显问题:

  • 接口响应时间变长
  • 下游服务故障会直接拖垮主流程
  • 系统耦合度高,难以扩展
  • 突发流量容易导致服务雪崩

引入消息队列后,主流程只需要把事件投递出去,后续动作交给消费者异步处理。这样可以实现:

  • 解耦:生产者不需要关心消费者内部实现
  • 异步化:缩短核心请求响应时间
  • 削峰填谷:缓冲突发流量
  • 广播通知:一个事件触发多个下游系统
  • 失败重试:通过消息重投提升系统韧性

二、技术选型对比

在开始编码前,先看一下三类常见消息系统的特点。

中间件 适用场景 优势 注意点
Kafka 日志流、埋点、订单流、数据管道 吞吐高、分区扩展强、持久化能力好 运维复杂度相对较高
RabbitMQ 业务消息、任务队列、延迟队列 路由能力强、协议成熟、功能丰富 大吞吐场景不如 Kafka
NATS 轻量级消息通信、云原生服务间事件分发 部署简单、延迟低、上手快 高级持久化能力需依赖 JetStream

如果你的核心诉求是高吞吐事件流,优先考虑 Kafka;如果是复杂业务路由和延迟任务,RabbitMQ 更常见;如果需要轻量、快速、云原生友好的消息通信,NATS 是非常合适的选择。

三、示例工程说明

为了便于演示,本文按多个独立样例展开。你可以创建一个目录 mq-demo,然后把下面的代码按文件保存。

建议使用如下目录结构:

mq-demo/
├── docker-compose.yml
├── go.mod
├── kafka/
│   ├── producer.go
│   └── consumer.go
├── rabbitmq/
│   ├── producer.go
│   ├── consumer.go
│   ├── delay_producer.go
│   └── delay_consumer.go
├── nats/
│   ├── pub.go
│   └── sub.go
├── eda/
│   └── order_event_demo.go
├── reliability/
│   └── idempotent_consumer.go
└── scheduler/
    └── cron_task.go

下面先准备统一依赖。

1. go.mod

module mq-demo

go 1.22

require (
	github.com/IBM/sarama v1.45.0
	github.com/nats-io/nats.go v1.39.1
	github.com/rabbitmq/amqp091-go v1.10.0
	github.com/robfig/cron/v3 v3.0.1
	modernc.org/sqlite v1.34.5
)

2. docker-compose.yml

为了让 Kafka、RabbitMQ、NATS 可以快速启动,我们使用如下容器配置。

version: "3.9"

services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.5.0
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
    ports:
      - "2181:2181"

  kafka:
    image: confluentinc/cp-kafka:7.5.0
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true"

  rabbitmq:
    image: rabbitmq:3.13-management
    ports:
      - "5672:5672"
      - "15672:15672"

  nats:
    image: nats:2.10-alpine
    ports:
      - "4222:4222"

启动命令:

docker compose up -d

四、Kafka 生产者与消费者实现

Kafka 的核心概念包括:

  • Topic:消息主题
  • Partition:分区,用于并行扩展
  • Producer:生产者,负责发送消息
  • Consumer Group:消费组,用于负载均衡消费
  • Offset:消息位点,用于标识消费进度

1. Kafka 生产者

下面使用 sarama 实现一个同步生产者,把订单创建事件发送到 orders 主题。

文件:kafka/producer.go

package main

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

	"github.com/IBM/sarama"
)

type OrderCreatedEvent struct {
	EventID   string    `json:"event_id"`
	OrderID   string    `json:"order_id"`
	UserID    int64     `json:"user_id"`
	Amount    float64   `json:"amount"`
	CreatedAt time.Time `json:"created_at"`
}

func main() {
	config := sarama.NewConfig()
	config.Version = sarama.V2_8_0_0
	config.Producer.RequiredAcks = sarama.WaitForAll
	config.Producer.Return.Successes = true
	config.Producer.Retry.Max = 3

	producer, err := sarama.NewSyncProducer([]string{"localhost:9092"}, config)
	if err != nil {
		log.Fatalf("create producer failed: %v", err)
	}
	defer producer.Close()

	event := OrderCreatedEvent{
		EventID:   "evt-kafka-1001",
		OrderID:   "order-1001",
		UserID:    20001,
		Amount:    99.90,
		CreatedAt: time.Now(),
	}

	body, err := json.Marshal(event)
	if err != nil {
		log.Fatalf("marshal event failed: %v", err)
	}

	msg := &sarama.ProducerMessage{
		Topic: "orders",
		Key:   sarama.StringEncoder(event.OrderID),
		Value: sarama.ByteEncoder(body),
		Headers: []sarama.RecordHeader{
			{Key: []byte("event_type"), Value: []byte("order_created")},
		},
	}

	partition, offset, err := producer.SendMessage(msg)
	if err != nil {
		log.Fatalf("send message failed: %v", err)
	}

	fmt.Printf("message sent successfully, partition=%d offset=%d\n", partition, offset)
}

运行:

go run kafka/producer.go

2. Kafka 消费者

下面实现一个消费组消费者,持续消费 orders 主题中的消息。

文件:kafka/consumer.go

package main

import (
	"context"
	"fmt"
	"log"
	"os"
	"os/signal"
	"sync"
	"syscall"

	"github.com/IBM/sarama"
)

type OrderConsumerGroupHandler struct{}

func (h *OrderConsumerGroupHandler) Setup(_ sarama.ConsumerGroupSession) error {
	fmt.Println("consumer group setup")
	return nil
}

func (h *OrderConsumerGroupHandler) Cleanup(_ sarama.ConsumerGroupSession) error {
	fmt.Println("consumer group cleanup")
	return nil
}

func (h *OrderConsumerGroupHandler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
	for msg := range claim.Messages() {
		fmt.Printf("received message topic=%s partition=%d offset=%d key=%s value=%s\n",
			msg.Topic, msg.Partition, msg.Offset, string(msg.Key), string(msg.Value))

		// 业务处理逻辑
		fmt.Println("processing order event...")

		// 标记消息已消费
		session.MarkMessage(msg, "processed")
	}
	return nil
}

func main() {
	config := sarama.NewConfig()
	config.Version = sarama.V2_8_0_0
	config.Consumer.Group.Rebalance.Strategy = sarama.NewBalanceStrategyRoundRobin()
	config.Consumer.Offsets.Initial = sarama.OffsetNewest

	group, err := sarama.NewConsumerGroup([]string{"localhost:9092"}, "order-service-group", config)
	if err != nil {
		log.Fatalf("create consumer group failed: %v", err)
	}
	defer group.Close()

	ctx, cancel := context.WithCancel(context.Background())
	wg := &sync.WaitGroup{}
	wg.Add(1)

	go func() {
		defer wg.Done()
		handler := &OrderConsumerGroupHandler{}
		for {
			if err := group.Consume(ctx, []string{"orders"}, handler); err != nil {
				log.Printf("consume failed: %v", err)
			}
			if ctx.Err() != nil {
				return
			}
		}
	}()

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

	cancel()
	wg.Wait()
	fmt.Println("consumer exited")
}

运行:

go run kafka/consumer.go

3. Kafka 使用要点

要点 说明
分区键设计 相同业务主键应尽量落到同一分区,方便保证局部顺序
消费组 同组内一个分区同一时刻只会被一个消费者实例消费
ACK 策略 生产端建议设置 WaitForAll 提高可靠性
Offset 提交 建议业务处理成功后再提交,避免消息丢失
批量发送 大吞吐场景可采用异步生产或批量发送提升性能

五、RabbitMQ 集成

RabbitMQ 更偏向业务消息处理。它的优势在于:

  • 支持多种交换机模型
  • 路由灵活
  • 易于构建任务队列
  • 在延迟队列、重试队列、死信队列方面工程实践成熟

1. RabbitMQ 生产者

文件:rabbitmq/producer.go

package main

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

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

type PaymentSucceededEvent struct {
	EventID   string    `json:"event_id"`
	OrderID   string    `json:"order_id"`
	Status    string    `json:"status"`
	CreatedAt time.Time `json:"created_at"`
}

func main() {
	conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
	if err != nil {
		log.Fatalf("connect rabbitmq failed: %v", err)
	}
	defer conn.Close()

	ch, err := conn.Channel()
	if err != nil {
		log.Fatalf("open channel failed: %v", err)
	}
	defer ch.Close()

	err = ch.ExchangeDeclare(
		"biz.events",
		"direct",
		true,
		false,
		false,
		false,
		nil,
	)
	if err != nil {
		log.Fatalf("declare exchange failed: %v", err)
	}

	event := PaymentSucceededEvent{
		EventID:   "evt-rabbit-1001",
		OrderID:   "order-1001",
		Status:    "SUCCESS",
		CreatedAt: time.Now(),
	}

	body, err := json.Marshal(event)
	if err != nil {
		log.Fatalf("marshal failed: %v", err)
	}

	err = ch.Publish(
		"biz.events",
		"payment.success",
		false,
		false,
		amqp.Publishing{
			ContentType:  "application/json",
			DeliveryMode: amqp.Persistent,
			Body:         body,
			MessageId:    event.EventID,
			Timestamp:    time.Now(),
		},
	)
	if err != nil {
		log.Fatalf("publish failed: %v", err)
	}

	log.Println("rabbitmq message published")
}

2. RabbitMQ 消费者

文件:rabbitmq/consumer.go

package main

import (
	"log"

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

func main() {
	conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
	if err != nil {
		log.Fatalf("connect rabbitmq failed: %v", err)
	}
	defer conn.Close()

	ch, err := conn.Channel()
	if err != nil {
		log.Fatalf("open channel failed: %v", err)
	}
	defer ch.Close()

	q, err := ch.QueueDeclare(
		"payment.success.queue",
		true,
		false,
		false,
		false,
		nil,
	)
	if err != nil {
		log.Fatalf("declare queue failed: %v", err)
	}

	err = ch.QueueBind(
		q.Name,
		"payment.success",
		"biz.events",
		false,
		nil,
	)
	if err != nil {
		log.Fatalf("bind queue failed: %v", err)
	}

	err = ch.Qos(1, 0, false)
	if err != nil {
		log.Fatalf("qos failed: %v", err)
	}

	msgs, err := ch.Consume(
		q.Name,
		"",
		false,
		false,
		false,
		false,
		nil,
	)
	if err != nil {
		log.Fatalf("consume failed: %v", err)
	}

	log.Println("waiting for rabbitmq messages...")
	forever := make(chan struct{})

	go func() {
		for msg := range msgs {
			log.Printf("received message: %s", string(msg.Body))

			// 模拟业务处理
			log.Println("handle payment success event")

			if err := msg.Ack(false); err != nil {
				log.Printf("ack failed: %v", err)
			}
		}
	}()

	<-forever
}

运行顺序:

go run rabbitmq/consumer.go
go run rabbitmq/producer.go

六、NATS 集成

NATS 的接口简洁,非常适合微服务之间的轻量级消息传递。在对实时性、部署简洁度要求较高的场景中,NATS 非常好用。

1. NATS 发布者

文件:nats/pub.go

package main

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

	"github.com/nats-io/nats.go"
)

type UserRegisteredEvent struct {
	EventID   string    `json:"event_id"`
	UserID    int64     `json:"user_id"`
	Email     string    `json:"email"`
	CreatedAt time.Time `json:"created_at"`
}

func main() {
	nc, err := nats.Connect(nats.DefaultURL)
	if err != nil {
		log.Fatalf("connect nats failed: %v", err)
	}
	defer nc.Close()

	event := UserRegisteredEvent{
		EventID:   "evt-nats-1001",
		UserID:    30001,
		Email:     "demo@example.com",
		CreatedAt: time.Now(),
	}

	body, err := json.Marshal(event)
	if err != nil {
		log.Fatalf("marshal failed: %v", err)
	}

	if err := nc.Publish("user.registered", body); err != nil {
		log.Fatalf("publish failed: %v", err)
	}

	if err := nc.Flush(); err != nil {
		log.Fatalf("flush failed: %v", err)
	}

	log.Println("nats message published")
}

2. NATS 订阅者

文件:nats/sub.go

package main

import (
	"log"

	"github.com/nats-io/nats.go"
)

func main() {
	nc, err := nats.Connect(nats.DefaultURL)
	if err != nil {
		log.Fatalf("connect nats failed: %v", err)
	}
	defer nc.Close()

	_, err = nc.Subscribe("user.registered", func(msg *nats.Msg) {
		log.Printf("received nats message: %s", string(msg.Data))
	})
	if err != nil {
		log.Fatalf("subscribe failed: %v", err)
	}

	log.Println("waiting for nats messages...")
	select {}
}

运行顺序:

go run nats/sub.go
go run nats/pub.go

3. RabbitMQ 与 NATS 的使用差异

维度 RabbitMQ NATS
路由能力 强,交换机模型丰富 相对简洁,以 Subject 为主
上手成本 中等
延迟表现 较好 很好
业务任务队列 很适合 也可做,但传统业务更多用 RabbitMQ
云原生友好度 较高 非常高

七、事件驱动架构

消息队列真正发挥价值的核心,不只是“发消息”和“收消息”,而是建立起事件驱动架构(EDA,Event-Driven Architecture)

1. 什么是事件驱动架构

事件驱动架构的核心思想是:

  • 业务动作发生后,不直接调用多个下游服务
  • 而是先发布一个“领域事件”
  • 由不同的订阅方根据事件独立处理自己的业务

例如,订单创建成功后,可以发出 OrderCreated 事件,随后:

  • 库存服务扣减库存
  • 营销服务发放优惠券
  • 通知服务发送短信
  • 数据平台更新报表

这样每个系统都可以独立演进,不会把主业务流程绑死在一条同步链路中。

2. 事件模型设计建议

设计事件时,推荐统一使用“事件信封”结构:

字段 说明
event_id 全局唯一事件 ID,用于追踪和幂等
event_type 事件类型,例如 order.created
source 事件来源服务
occurred_at 事件发生时间
data 业务载荷

下面给出一个简化版事件发布示例。

文件:eda/order_event_demo.go

package main

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

type Event struct {
	EventID    string      `json:"event_id"`
	EventType  string      `json:"event_type"`
	Source     string      `json:"source"`
	OccurredAt time.Time   `json:"occurred_at"`
	Data       interface{} `json:"data"`
}

type OrderCreatedData struct {
	OrderID string  `json:"order_id"`
	UserID  int64   `json:"user_id"`
	Amount  float64 `json:"amount"`
}

type EventPublisher interface {
	Publish(topic string, event Event) error
}

type ConsolePublisher struct{}

func (c *ConsolePublisher) Publish(topic string, event Event) error {
	body, err := json.MarshalIndent(event, "", "  ")
	if err != nil {
		return err
	}
	fmt.Printf("publish to topic=%s\n%s\n", topic, string(body))
	return nil
}

type OrderService struct {
	publisher EventPublisher
}

func NewOrderService(publisher EventPublisher) *OrderService {
	return &OrderService{publisher: publisher}
}

func (s *OrderService) CreateOrder(orderID string, userID int64, amount float64) error {
	log.Printf("create order success: orderID=%s userID=%d amount=%.2f", orderID, userID, amount)

	event := Event{
		EventID:    "evt-eda-1001",
		EventType:  "order.created",
		Source:     "order-service",
		OccurredAt: time.Now(),
		Data: OrderCreatedData{
			OrderID: orderID,
			UserID:  userID,
			Amount:  amount,
		},
	}

	return s.publisher.Publish("order.created", event)
}

func main() {
	publisher := &ConsolePublisher{}
	service := NewOrderService(publisher)

	if err := service.CreateOrder("order-2001", 40001, 199.00); err != nil {
		log.Fatalf("create order failed: %v", err)
	}
}

这个例子中虽然只是打印控制台,但在实际项目里,ConsolePublisher 可以很容易替换为 Kafka、RabbitMQ 或 NATS 实现。这样业务层只依赖 EventPublisher 接口,而不直接绑死具体消息系统。

3. 事件驱动架构的落地建议

  • 事件命名统一,例如 order.createdpayment.succeeded
  • 事件结构版本化,避免字段变化导致消费者崩溃
  • 消费者必须可重复消费,避免因为重试造成脏数据
  • 不要在事件里塞入过多冗余数据,避免事件膨胀
  • 核心链路建议保留链路追踪字段,如 trace_id

八、消息可靠性保证与幂等设计

消息系统中最常见的问题,不是“消息发不出去”,而是:

  • 发出去了但消费者处理失败
  • 消费者处理成功了但 ACK 丢失,导致重复投递
  • 生产者重试造成重复消息
  • 服务重启时消息状态不一致

因此,消息可靠性幂等设计必须一起考虑。

1. 常见投递语义

语义 含义 特点
At Most Once 最多一次 可能丢消息,但不会重复
At Least Once 至少一次 不会轻易丢消息,但可能重复
Exactly Once 恰好一次 实现成本高,通常需中间件与业务联合保证

在大多数业务系统中,现实可行的方案是:

使用至少一次投递 + 消费端幂等处理。

2. 生产端可靠性建议

生产端常见做法包括:

  • 开启发送确认机制
  • 失败自动重试
  • 使用持久化消息
  • 关键业务采用本地消息表 / Outbox 模式

其中,Outbox 模式非常值得掌握:

  1. 本地事务中同时写业务数据和消息表
  2. 后台任务扫描消息表
  3. 把消息可靠投递到 MQ
  4. 投递成功后更新消息状态

这样可以避免“数据库写成功但 MQ 发送失败”的双写一致性问题。

3. 消费端幂等设计

消费者需要能识别“这条消息是否已经处理过”。常见做法包括:

  • 基于消息唯一 ID 去重
  • 数据库唯一索引防重
  • Redis SETNX 防重
  • 业务状态机校验,例如订单已支付则不重复扣款

下面给出一个可运行的幂等消费示例。它使用 SQLite 记录已处理消息 ID,重复消息不会重复执行业务逻辑。

文件:reliability/idempotent_consumer.go

package main

import (
	"database/sql"
	"fmt"
	"log"
	"time"

	_ "modernc.org/sqlite"
)

type Message struct {
	EventID string
	OrderID string
	Amount  float64
}

func main() {
	db, err := sql.Open("sqlite", "file:idempotent_demo.db?cache=shared&mode=rwc")
	if err != nil {
		log.Fatalf("open db failed: %v", err)
	}
	defer db.Close()

	if err := initDB(db); err != nil {
		log.Fatalf("init db failed: %v", err)
	}

	messages := []Message{
		{EventID: "evt-5001", OrderID: "order-5001", Amount: 88.80},
		{EventID: "evt-5001", OrderID: "order-5001", Amount: 88.80},
		{EventID: "evt-5002", OrderID: "order-5002", Amount: 66.60},
	}

	for _, msg := range messages {
		if err := handleMessage(db, msg); err != nil {
			log.Printf("handle message failed: %v", err)
		}
	}
}

func initDB(db *sql.DB) error {
	createTableSQL := `
CREATE TABLE IF NOT EXISTS consumed_messages (
    event_id TEXT PRIMARY KEY,
    consumed_at DATETIME NOT NULL
);
`
	_, err := db.Exec(createTableSQL)
	return err
}

func handleMessage(db *sql.DB, msg Message) error {
	tx, err := db.Begin()
	if err != nil {
		return err
	}
	defer tx.Rollback()

	res, err := tx.Exec(
		"INSERT OR IGNORE INTO consumed_messages(event_id, consumed_at) VALUES (?, ?)",
		msg.EventID,
		time.Now(),
	)
	if err != nil {
		return err
	}

	rowsAffected, err := res.RowsAffected()
	if err != nil {
		return err
	}

	if rowsAffected == 0 {
		fmt.Printf("skip duplicated message: eventID=%s orderID=%s\n", msg.EventID, msg.OrderID)
		return tx.Commit()
	}

	// 这里模拟真正的业务逻辑,例如发券、记账、更新订单状态等。
	fmt.Printf("process message: eventID=%s orderID=%s amount=%.2f\n", msg.EventID, msg.OrderID, msg.Amount)

	return tx.Commit()
}

运行:

go run reliability/idempotent_consumer.go

预期输出类似:

process message: eventID=evt-5001 orderID=order-5001 amount=88.80
skip duplicated message: eventID=evt-5001 orderID=order-5001
process message: eventID=evt-5002 orderID=order-5002 amount=66.60

4. 幂等设计的工程经验

场景 推荐方案
订单支付回调 以支付流水号或回调事件 ID 做唯一约束
发券 / 发红包 以用户 ID + 活动 ID + 业务单号做幂等键
库存扣减 结合订单号与状态机,避免重复扣减
短信 / 邮件发送 记录发送流水,防止重复通知

5. 重试与死信队列

仅有幂等还不够,生产环境还应配合:

  • 有限重试:临时性故障可重试
  • 退避策略:如指数退避,防止雪崩
  • 死信队列:多次失败后转入死信队列,人工或系统补偿
  • 告警监控:堆积、失败率、重试次数应纳入监控

九、延迟队列与定时任务

很多业务都不是“消息一到就立刻处理”,而是希望:

  • 30 分钟未支付自动取消订单
  • 10 分钟后发送提醒通知
  • 每天凌晨对账
  • 每小时扫描超时任务

这类需求通常通过延迟队列定时任务实现。

1. RabbitMQ TTL + 死信队列实现延迟消息

RabbitMQ 常见的延迟方案是:

  1. 先把消息投递到延迟队列
  2. 为队列或消息设置 TTL
  3. TTL 到期后转发到死信交换机
  4. 最终消费者从目标队列中消费

下面给出完整样例。

文件:rabbitmq/delay_consumer.go

package main

import (
	"log"

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

func main() {
	conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
	if err != nil {
		log.Fatalf("connect rabbitmq failed: %v", err)
	}
	defer conn.Close()

	ch, err := conn.Channel()
	if err != nil {
		log.Fatalf("open channel failed: %v", err)
	}
	defer ch.Close()

	if err := declareDelayInfra(ch); err != nil {
		log.Fatalf("declare delay infra failed: %v", err)
	}

	msgs, err := ch.Consume(
		"order.release.queue",
		"",
		false,
		false,
		false,
		false,
		nil,
	)
	if err != nil {
		log.Fatalf("consume failed: %v", err)
	}

	log.Println("waiting delayed messages...")
	forever := make(chan struct{})

	go func() {
		for msg := range msgs {
			log.Printf("received delayed message: %s", string(msg.Body))
			log.Println("execute order close or inventory rollback")
			msg.Ack(false)
		}
	}()

	<-forever
}

func declareDelayInfra(ch *amqp.Channel) error {
	if err := ch.ExchangeDeclare("order.dlx", "direct", true, false, false, false, nil); err != nil {
		return err
	}

	_, err := ch.QueueDeclare(
		"order.release.queue",
		true,
		false,
		false,
		false,
		nil,
	)
	if err != nil {
		return err
	}

	if err := ch.QueueBind("order.release.queue", "order.release", "order.dlx", false, nil); err != nil {
		return err
	}

	args := amqp.Table{
		"x-dead-letter-exchange":    "order.dlx",
		"x-dead-letter-routing-key": "order.release",
	}

	_, err = ch.QueueDeclare(
		"order.delay.queue",
		true,
		false,
		false,
		false,
		args,
	)
	return err
}

文件:rabbitmq/delay_producer.go

package main

import (
	"log"

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

func main() {
	conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
	if err != nil {
		log.Fatalf("connect rabbitmq failed: %v", err)
	}
	defer conn.Close()

	ch, err := conn.Channel()
	if err != nil {
		log.Fatalf("open channel failed: %v", err)
	}
	defer ch.Close()

	body := []byte(`{"order_id":"order-9001","action":"auto_cancel"}`)

	err = ch.Publish(
		"",
		"order.delay.queue",
		false,
		false,
		amqp.Publishing{
			ContentType:  "application/json",
			DeliveryMode: amqp.Persistent,
			Body:         body,
			Expiration:   "10000",
		},
	)
	if err != nil {
		log.Fatalf("publish delay message failed: %v", err)
	}

	log.Println("delay message published, will be delivered after 10 seconds")
}

运行顺序:

go run rabbitmq/delay_consumer.go
go run rabbitmq/delay_producer.go

10 秒后,消费者会从 order.release.queue 收到消息。

2. Go 定时任务实现

不是所有延迟需求都要依赖 MQ。对于周期性任务,例如“每天 1 点跑报表”,直接使用调度器会更直观。

下面使用 robfig/cron 实现一个简单的定时任务。

文件:scheduler/cron_task.go

package main

import (
	"log"
	"time"

	"github.com/robfig/cron/v3"
)

func main() {
	c := cron.New(cron.WithSeconds())

	_, err := c.AddFunc("0 */15 * * * *", func() {
		log.Printf("scan timeout orders at %s", time.Now().Format(time.DateTime))
	})
	if err != nil {
		log.Fatalf("add cron task failed: %v", err)
	}

	c.Start()
	log.Println("cron started, execute every 15 minutes")

	select {}
}

运行:

go run scheduler/cron_task.go

3. 延迟队列与定时任务如何选

场景 推荐方案
与某个业务事件强绑定,例如下单后 30 分钟取消 延迟队列
固定周期执行,如每小时扫描一次超时数据 定时任务
需要与消息消费链路统一治理 延迟队列
任务量固定、逻辑简单 定时任务

十、生产实践建议

把消息队列真正用好,往往不在于 API 会不会调用,而在于工程细节是否扎实。以下是几个非常关键的建议。

1. 统一消息模型

建议团队统一消息结构,例如:

{
  "event_id": "evt-10001",
  "event_type": "order.created",
  "trace_id": "trace-abc-001",
  "source": "order-service",
  "occurred_at": "2026-06-08T09:00:00+08:00",
  "data": {
    "order_id": "order-10001",
    "user_id": 12345,
    "amount": 88.8
  }
}

统一结构有几个好处:

  • 便于追踪与排障
  • 便于消费者做通用解析
  • 便于做链路监控和审计
  • 便于后续版本演进

2. 消费失败不要直接丢弃

正确做法通常是:

  1. 记录错误日志
  2. 做有限次数重试
  3. 超过阈值进入死信队列
  4. 配合告警和人工补偿

如果消费者异常时直接 Ack,问题会被悄悄吞掉;如果一直不 Ack 又不做治理,也可能导致队列阻塞。

3. 控制消息粒度

消息不要设计得过大。一般建议:

  • 只传必要字段
  • 大对象通过对象存储或数据库引用
  • 避免把整张订单快照无限膨胀地塞进消息体

4. 保证消费者无状态或弱状态

消费实例最好可以随时横向扩缩容,因此:

  • 不要把重要状态只保存在进程内存中
  • 去重状态应落 Redis / DB
  • 消费逻辑尽量可重入、可恢复

5. 监控指标必须齐全

至少应覆盖以下维度:

指标 说明
生产成功率 监控发送失败和重试情况
消费成功率 监控业务处理异常
消息堆积量 反映消费者是否跟不上
消费延迟 从产生到消费的时间差
死信数量 反映异常消息规模

十一、小结

本文围绕 Golang 在消息队列与异步处理中的核心实践,系统讲解了以下内容:

  • 使用 Kafka 实现生产者与消费者,理解 Topic、Partition、Consumer Group 等关键概念
  • 使用 RabbitMQ 处理业务消息,使用 NATS 实现轻量级服务间事件通信
  • 通过事件驱动架构解耦业务流程,让系统更易扩展
  • 通过 ACK、重试、幂等、Outbox 等手段保障消息可靠性
  • 通过延迟队列和定时任务实现超时关闭、延迟通知、周期扫描等业务需求

当你在项目中引入消息队列时,真正决定系统质量的,往往不是选了哪一种中间件,而是是否把以下问题处理扎实:

  • 消息是否会丢
  • 消息是否会重复
  • 消费失败如何补偿
  • 延迟任务是否可控
  • 整个链路是否可观测

只要把这些基础能力建立起来,消息驱动系统就能成为高并发后端架构中非常稳定的一块基石。


📝 版权声明:本文为原创技术博客,转载请注明出处。

如文章中存在错误或不准确之处,欢迎在评论区指正,感谢您的阅读与支持!

上一篇

Golang项目实践:4.4 微服务架构实战之构建可观测、可扩展的服务体系

下一篇

Golang项目实践:4.8 综合项目之分布式短链接服务