消息队列
Broker 无关消息队列,RabbitMQ、Kafka 等后端
在线 Playground
在浏览器中直接体验本页相关 API,无需本地安装 Go 环境。
另见专题文档:消息队列完整文档。以下为
common/mq包 README 全文。
mq
Broker-agnostic message queue abstraction with pluggable backends.
Packages
| Package | Description |
|---|---|
mq/ | Core interfaces: Message, Delivery, Producer, Consumer, Broker, Handler, Middleware |
mq/factory/ | Unified factory — NewBroker(BrokerConfig) dispatches to all backends |
mq/rabbitmq/ | RabbitMQ backend (amqp091-go) |
mq/kafka/ | Kafka backend (segmentio/kafka-go) |
mq/rocketmq/ | RocketMQ backend (apache/rocketmq-client-go/v2) |
mq/activemq/ | ActiveMQ backend via STOMP (go-stomp/stomp) |
mq/redisstream/ | Redis Streams backend (go-redis) |
Quick start via factory
import (
"github.com/LingByte/ling-base/common/mq"
"github.com/LingByte/ling-base/common/mq/factory"
"github.com/LingByte/ling-base/common/mq/rabbitmq"
)
broker, err := factory.NewBroker(mq.BrokerConfig{
Type: mq.BrokerRabbitMQ,
RabbitMQConfig: &rabbitmq.Config{URL: "amqp://guest:guest@localhost:5672/"},
})
if err != nil { panic(err) }
defer broker.Close()
broker.Connect()
// Produce
producer, _ := broker.Producer("events", mq.PublishOptions{Persistent: true})
producer.Publish(ctx, &mq.Message{
RoutingKey: "events.user.created",
ContentType: "application/json",
Body: []byte(`{"user_id":123}`),
})
// Consume with middleware
consumer, _ := broker.Consumer("events.queue", mq.ConsumeOptions{
Handler: func(ctx context.Context, d mq.Delivery) error {
log.Println(string(d.Body()))
return d.Ack()
},
Middleware: []mq.Middleware{
mq.RecoverMiddleware(logger),
mq.MetricsMiddleware(metrics),
mq.RetryMiddleware(mq.RetryConfig{MaxAttempts: 3, Backoff: time.Second}),
},
QosPrefetchCount: 10,
Concurrency: 4,
})
consumer.Start(ctx)Switching backends
// Kafka
broker, _ := factory.NewBroker(mq.BrokerConfig{
Type: mq.BrokerKafka,
KafkaConfig: &kafka.Config{Brokers: []string{"localhost:9092"}, GroupID: "my-group"},
})
// RocketMQ
broker, _ := factory.NewBroker(mq.BrokerConfig{
Type: mq.BrokerRocketMQ,
RocketMQConfig: &rocketmq.Config{NameServer: []string{"127.0.0.1:9876"}},
})
// ActiveMQ (STOMP)
broker, _ := factory.NewBroker(mq.BrokerConfig{
Type: mq.BrokerActiveMQ,
ActiveMQConfig: &activemq.Config{Addr: "localhost:61613"},
})
// Redis Streams
broker, _ := factory.NewBroker(mq.BrokerConfig{
Type: mq.BrokerRedisStream,
RedisStreamConfig: &redisstream.Config{Addr: "localhost:6379", Group: "my-group"},
})Backend capabilities
| Backend | Persistent | Ordered | Transaction | ConsumerGroup | PubSub |
|---|---|---|---|---|---|
| RabbitMQ | ✅ | ✅ | ✅ | ❌ | ✅ |
| Kafka | ✅ | ✅ | ✅ | ✅ | ✅ |
| RocketMQ | ✅ | ✅ | ✅ | ✅ | ✅ |
| ActiveMQ | ✅ | ✅ | ❌ | ❌ | ✅ |
| Redis Streams | ✅ | ✅ | ❌ | ✅ | ❌ |
Middleware
| Middleware | Description |
|---|---|
RecoverMiddleware | Recovers from handler panics |
LoggingMiddleware | Logs each delivery before/after |
MetricsMiddleware | Records consume/error/redelivered metrics |
RetryMiddleware | Retries failed deliveries with backoff |
DeadLetterMiddleware | Routes failures to a dead-letter handler |
Chain | Composes multiple middleware |