在微服务、订单系统、支付系统、通知系统等场景中,**消息队列(Message Queue)**是实现系统解耦、削峰填谷、异步处理和事件驱动的重要基础设施。对于 Golang 开发者来说,掌握 Kafka、RabbitMQ、NATS 等常见消息中间件的使用方式,以及消息可靠性、幂等、延迟任务等工程实践,是构建高可用后端系统的关键能力。
本文将围绕以下几个方面展开:
- Kafka 生产者与消费者实现
- RabbitMQ / NATS 集成
- 事件驱动架构
- 消息可靠性保证与幂等设计
- 延迟队列与定时任务
文章中的示例均采用 Go 编写,并尽量提供可直接运行的代码。
一、为什么需要消息队列
在同步调用模型里,一个请求往往要串行完成多个动作,例如:
- 创建订单
- 扣减库存
- 发送短信
- 写审计日志
- 触发推荐系统更新
如果这些动作全部放在一个同步链路中,会带来几个明显问题:
- 接口响应时间变长
- 下游服务故障会直接拖垮主流程
- 系统耦合度高,难以扩展
- 突发流量容易导致服务雪崩
引入消息队列后,主流程只需要把事件投递出去,后续动作交给消费者异步处理。这样可以实现:
- 解耦:生产者不需要关心消费者内部实现
- 异步化:缩短核心请求响应时间
- 削峰填谷:缓冲突发流量
- 广播通知:一个事件触发多个下游系统
- 失败重试:通过消息重投提升系统韧性
二、技术选型对比
在开始编码前,先看一下三类常见消息系统的特点。
| 中间件 | 适用场景 | 优势 | 注意点 |
|---|---|---|---|
| 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.created、payment.succeeded - 事件结构版本化,避免字段变化导致消费者崩溃
- 消费者必须可重复消费,避免因为重试造成脏数据
- 不要在事件里塞入过多冗余数据,避免事件膨胀
- 核心链路建议保留链路追踪字段,如
trace_id
八、消息可靠性保证与幂等设计
消息系统中最常见的问题,不是“消息发不出去”,而是:
- 发出去了但消费者处理失败
- 消费者处理成功了但 ACK 丢失,导致重复投递
- 生产者重试造成重复消息
- 服务重启时消息状态不一致
因此,消息可靠性和幂等设计必须一起考虑。
1. 常见投递语义
| 语义 | 含义 | 特点 |
|---|---|---|
| At Most Once | 最多一次 | 可能丢消息,但不会重复 |
| At Least Once | 至少一次 | 不会轻易丢消息,但可能重复 |
| Exactly Once | 恰好一次 | 实现成本高,通常需中间件与业务联合保证 |
在大多数业务系统中,现实可行的方案是:
使用至少一次投递 + 消费端幂等处理。
2. 生产端可靠性建议
生产端常见做法包括:
- 开启发送确认机制
- 失败自动重试
- 使用持久化消息
- 关键业务采用本地消息表 / Outbox 模式
其中,Outbox 模式非常值得掌握:
- 本地事务中同时写业务数据和消息表
- 后台任务扫描消息表
- 把消息可靠投递到 MQ
- 投递成功后更新消息状态
这样可以避免“数据库写成功但 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 常见的延迟方案是:
- 先把消息投递到延迟队列
- 为队列或消息设置 TTL
- TTL 到期后转发到死信交换机
- 最终消费者从目标队列中消费
下面给出完整样例。
文件: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. 消费失败不要直接丢弃
正确做法通常是:
- 记录错误日志
- 做有限次数重试
- 超过阈值进入死信队列
- 配合告警和人工补偿
如果消费者异常时直接 Ack,问题会被悄悄吞掉;如果一直不 Ack 又不做治理,也可能导致队列阻塞。
3. 控制消息粒度
消息不要设计得过大。一般建议:
- 只传必要字段
- 大对象通过对象存储或数据库引用
- 避免把整张订单快照无限膨胀地塞进消息体
4. 保证消费者无状态或弱状态
消费实例最好可以随时横向扩缩容,因此:
- 不要把重要状态只保存在进程内存中
- 去重状态应落 Redis / DB
- 消费逻辑尽量可重入、可恢复
5. 监控指标必须齐全
至少应覆盖以下维度:
| 指标 | 说明 |
|---|---|
| 生产成功率 | 监控发送失败和重试情况 |
| 消费成功率 | 监控业务处理异常 |
| 消息堆积量 | 反映消费者是否跟不上 |
| 消费延迟 | 从产生到消费的时间差 |
| 死信数量 | 反映异常消息规模 |
十一、小结
本文围绕 Golang 在消息队列与异步处理中的核心实践,系统讲解了以下内容:
- 使用 Kafka 实现生产者与消费者,理解 Topic、Partition、Consumer Group 等关键概念
- 使用 RabbitMQ 处理业务消息,使用 NATS 实现轻量级服务间事件通信
- 通过事件驱动架构解耦业务流程,让系统更易扩展
- 通过 ACK、重试、幂等、Outbox 等手段保障消息可靠性
- 通过延迟队列和定时任务实现超时关闭、延迟通知、周期扫描等业务需求
当你在项目中引入消息队列时,真正决定系统质量的,往往不是选了哪一种中间件,而是是否把以下问题处理扎实:
- 消息是否会丢
- 消息是否会重复
- 消费失败如何补偿
- 延迟任务是否可控
- 整个链路是否可观测
只要把这些基础能力建立起来,消息驱动系统就能成为高并发后端架构中非常稳定的一块基石。
📝 版权声明:本文为原创技术博客,转载请注明出处。
如文章中存在错误或不准确之处,欢迎在评论区指正,感谢您的阅读与支持!