ling-baseling-base

消息队列

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) Handler

RabbitMQ 完整示例

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

中间件

中间件作用
RecoverMiddlewarepanic 恢复,避免消费者崩溃
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.ErrClosedBroker 已关闭
mq.ErrNotConnected未 Connect
mq.ErrAlreadyRunningConsumer 重复 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 完成(视后端实现)

On this page