消息队列
RabbitMQ/Kafka/RocketMQ/ActiveMQ/Redis Streams 统一 Broker
在线 Playground
在浏览器中直接体验本页相关 API,无需本地安装 Go 环境。
消息队列 (mq)
common/mq 提供 Broker 无关的消息队列抽象,通过 factory.NewBroker 切换后端,业务代码(Producer / Consumer / Middleware)保持不变。
架构
common/mq/
├── mq.go # Message、Delivery、Producer、Consumer、Broker、Middleware
├── factory/ # NewBroker(BrokerConfig) 统一工厂
├── rabbitmq/ # RabbitMQ (amqp091-go)
├── kafka/ # Kafka (segmentio/kafka-go)
├── rocketmq/ # RocketMQ
├── activemq/ # ActiveMQ (STOMP)
└── redisstream/ # Redis Streams安装
go get github.com/LingByte/ling-base/common/mq
go get github.com/LingByte/ling-base/common/mq/factory
go get github.com/LingByte/ling-base/common/mq/rabbitmq核心类型
Message(发布)
type Message struct {
ID string
Exchange string // RabbitMQ exchange / Kafka topic
RoutingKey string
Headers map[string]any
Body []byte
ContentType string // 如 application/json
ContentEncoding string
Priority uint8
CorrelationID string // RPC 关联 ID
ReplyTo string
Expiration time.Duration
Timestamp time.Time
Type string
DeliveryMode DeliveryMode // Persistent / Transient
}Delivery(消费)
type Delivery interface {
Body() []byte
Headers() map[string]any
Ack() error
Nack(requeue bool) error
Reject(requeue bool) error
}Handler 与 Middleware
type Handler func(ctx context.Context, d Delivery) error
type Middleware func(Handler) HandlerRabbitMQ 完整示例
package main
import (
"context"
"log"
"time"
"github.com/LingByte/ling-base/common/mq"
"github.com/LingByte/ling-base/common/mq/factory"
"github.com/LingByte/ling-base/common/mq/rabbitmq"
)
func main() {
broker, err := factory.NewBroker(mq.BrokerConfig{
Type: mq.BrokerRabbitMQ,
RabbitMQConfig: &rabbitmq.Config{
URL: "amqp://guest:guest@localhost:5672/",
},
})
if err != nil {
log.Fatal(err)
}
defer broker.Close()
if err := broker.Connect(); err != nil {
log.Fatal(err)
}
ctx := context.Background()
// 生产者
producer, err := broker.Producer("events", mq.PublishOptions{
Persistent: true,
})
if err != nil {
log.Fatal(err)
}
err = producer.Publish(ctx, &mq.Message{
RoutingKey: "events.user.created",
ContentType: "application/json",
Body: []byte(`{"user_id":123,"email":"a@b.com"}`),
Headers: map[string]any{
"source": "api",
},
})
if err != nil {
log.Fatal(err)
}
// 消费者
consumer, err := broker.Consumer("events.queue", mq.ConsumeOptions{
Handler: handleDelivery,
Middleware: []mq.Middleware{
mq.RecoverMiddleware(log.Printf),
mq.RetryMiddleware(mq.RetryConfig{
MaxAttempts: 3,
Backoff: time.Second,
}),
mq.MetricsMiddleware(myMetrics),
},
QosPrefetchCount: 10,
Concurrency: 4,
})
if err != nil {
log.Fatal(err)
}
if err := consumer.Start(ctx); err != nil {
log.Fatal(err)
}
<-ctx.Done()
}
func handleDelivery(ctx context.Context, d mq.Delivery) error {
log.Printf("got message: %s", d.Body())
// 处理成功必须 Ack
return d.Ack()
}切换后端
Kafka
import "github.com/LingByte/ling-base/common/mq/kafka"
broker, _ := factory.NewBroker(mq.BrokerConfig{
Type: mq.BrokerKafka,
KafkaConfig: &kafka.Config{
Brokers: []string{"localhost:9092"},
GroupID: "my-service",
},
})RocketMQ
import "github.com/LingByte/ling-base/common/mq/rocketmq"
broker, _ := factory.NewBroker(mq.BrokerConfig{
Type: mq.BrokerRocketMQ,
RocketMQConfig: &rocketmq.Config{
NameServer: []string{"127.0.0.1:9876"},
},
})ActiveMQ (STOMP)
import "github.com/LingByte/ling-base/common/mq/activemq"
broker, _ := factory.NewBroker(mq.BrokerConfig{
Type: mq.BrokerActiveMQ,
ActiveMQConfig: &activemq.Config{Addr: "localhost:61613"},
})Redis Streams
import "github.com/LingByte/ling-base/common/mq/redisstream"
broker, _ := factory.NewBroker(mq.BrokerConfig{
Type: mq.BrokerRedisStream,
RedisStreamConfig: &redisstream.Config{
Addr: "localhost:6379",
Group: "my-group",
},
})后端能力对照
| 后端 | 持久化 | 有序 | 事务 | 消费组 | Pub/Sub |
|---|---|---|---|---|---|
| RabbitMQ | ✅ | ✅ | ✅ | ❌ | ✅ |
| Kafka | ✅ | ✅ | ✅ | ✅ | ✅ |
| RocketMQ | ✅ | ✅ | ✅ | ✅ | ✅ |
| ActiveMQ | ✅ | ✅ | ❌ | ❌ | ✅ |
| Redis Streams | ✅ | ✅ | ❌ | ✅ | ❌ |
中间件
| 中间件 | 作用 |
|---|---|
RecoverMiddleware | panic 恢复,避免消费者崩溃 |
LoggingMiddleware | 记录每条消息处理前后日志 |
MetricsMiddleware | 消费数、错误数、重投递指标 |
RetryMiddleware | 失败自动重试 + 退避 |
DeadLetterMiddleware | 超限失败转入死信处理 |
Chain | 组合多个中间件 |
handler := mq.Chain(
mq.RecoverMiddleware(log.Printf),
mq.LoggingMiddleware(log.Printf),
mq.RetryMiddleware(mq.RetryConfig{MaxAttempts: 5, Backoff: 2 * time.Second}),
)(myHandler)错误码
| 错误 | 含义 |
|---|---|
mq.ErrClosed | Broker 已关闭 |
mq.ErrNotConnected | 未 Connect |
mq.ErrAlreadyRunning | Consumer 重复 Start |
mq.ErrNoHandler | 未配置 Handler |
RPC 模式(Request/Reply)
// 发布时设置 CorrelationID 与 ReplyTo
_ = producer.Publish(ctx, &mq.Message{
Body: requestJSON,
CorrelationID: corrID,
ReplyTo: "rpc.reply.queue",
})
// 消费者在 ReplyTo 队列回复测试与运维
- 开发环境推荐 RabbitMQ +
docker compose - 生产环境为 Consumer 配置
RecoverMiddleware+RetryMiddleware+ 死信队列 broker.Close()会等待进行中的 ACK 完成(视后端实现)