ling-baseling-base

消息队列

Broker 无关消息队列,RabbitMQ、Kafka 等后端

在线 Playground

在浏览器中直接体验本页相关 API,无需本地安装 Go 环境。

另见专题文档:消息队列完整文档。以下为 common/mq 包 README 全文。

mq

Broker-agnostic message queue abstraction with pluggable backends.

Packages

PackageDescription
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

BackendPersistentOrderedTransactionConsumerGroupPubSub
RabbitMQ
Kafka
RocketMQ
ActiveMQ
Redis Streams

Middleware

MiddlewareDescription
RecoverMiddlewareRecovers from handler panics
LoggingMiddlewareLogs each delivery before/after
MetricsMiddlewareRecords consume/error/redelivered metrics
RetryMiddlewareRetries failed deliveries with backoff
DeadLetterMiddlewareRoutes failures to a dead-letter handler
ChainComposes multiple middleware

On this page