121 lines
3.2 KiB
Go
121 lines
3.2 KiB
Go
// Package mq broker 无关的消息队列抽象层入口。
|
|
//
|
|
// 启动流程:
|
|
//
|
|
// 1. 各服务 main.go 调用 mq.Init(cfg) — 内部根据 cfg.MQDriver 选择 adapter 实现
|
|
// 2. 业务侧通过 adapter.Get().TaskProducer()/EventProducer() 发送
|
|
// 3. 业务侧通过 adapter.Get().TaskConsumer()/EventConsumer() 注册 handler
|
|
// 4. main.go 调 mq.Close() 关闭连接
|
|
package mq
|
|
|
|
import (
|
|
"fmt"
|
|
"sync"
|
|
|
|
"github.com/topfans/backend/pkg/mq/adapter"
|
|
asynqAdapter "github.com/topfans/backend/pkg/mq/asynq"
|
|
streamsAdapter "github.com/topfans/backend/pkg/mq/streams"
|
|
)
|
|
|
|
var (
|
|
initOnce sync.Once
|
|
cfg Config
|
|
)
|
|
|
|
// Init 初始化 MQ adapter。
|
|
//
|
|
// 根据 cfg.MQDriver 选择实现:
|
|
//
|
|
// - "redis" (默认): asynq 处理任务 + Redis Streams 处理事件
|
|
// - "rabbitmq": 预留,本期不实现
|
|
//
|
|
// 重入安全,重复调用只生效一次。
|
|
func Init(config Config) error {
|
|
var initErr error
|
|
initOnce.Do(func() {
|
|
cfg = config
|
|
switch config.MQDriver {
|
|
case DriverRedis, "":
|
|
a, err := buildRedisAdapter(config)
|
|
if err != nil {
|
|
initErr = fmt.Errorf("mq: build redis adapter: %w", err)
|
|
return
|
|
}
|
|
adapter.Set(a)
|
|
case DriverRabbitMQ:
|
|
initErr = fmt.Errorf("mq: rabbitmq driver is reserved, not implemented yet")
|
|
default:
|
|
initErr = fmt.Errorf("mq: unknown driver %q", config.MQDriver)
|
|
return
|
|
}
|
|
})
|
|
return initErr
|
|
}
|
|
|
|
// Close 关闭底层 broker 连接。
|
|
func Close() error {
|
|
a := adapter.Get()
|
|
return a.Close()
|
|
}
|
|
|
|
// buildRedisAdapter 构造 redis 后端 asynq + redis-streams 双 adapter 并组装为单 Adapter。
|
|
func buildRedisAdapter(cfg Config) (adapter.Adapter, error) {
|
|
asynqA, err := asynqAdapter.New(cfg.Asynq)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("asynq adapter: %w", err)
|
|
}
|
|
streamsA, err := streamsAdapter.New(streamsAdapter.Config{
|
|
RedisAddr: cfg.Streams.RedisAddr,
|
|
RedisDB: cfg.Streams.RedisDB,
|
|
Password: cfg.Streams.Password,
|
|
ConsumerGroup: cfg.Streams.ConsumerGroup,
|
|
MaxLen: cfg.Streams.MaxLen,
|
|
ReadBlockMS: cfg.Streams.ReadBlockMS,
|
|
})
|
|
if err != nil {
|
|
_ = asynqA.Close()
|
|
return nil, fmt.Errorf("streams adapter: %w", err)
|
|
}
|
|
return &compositeAdapter{
|
|
taskAdapter: asynqA,
|
|
eventAdapter: streamsA,
|
|
}, nil
|
|
}
|
|
|
|
// compositeAdapter 把 asynq (TaskProducer/TaskConsumer) 和 streams (EventProducer/EventConsumer)
|
|
// 组装成单一 Adapter 实现。
|
|
//
|
|
// 设计上 asynq 提供 TaskProducer/TaskConsumer 的实现,streams 提供 EventProducer/EventConsumer。
|
|
// 它们各自负责自己的原语,不需要互相调用,因此合成关系简单。
|
|
type compositeAdapter struct {
|
|
taskAdapter *asynqAdapter.Adapter
|
|
eventAdapter *streamsAdapter.Adapter
|
|
}
|
|
|
|
func (c *compositeAdapter) TaskProducer() adapter.TaskProducer {
|
|
return c.taskAdapter
|
|
}
|
|
|
|
func (c *compositeAdapter) TaskConsumer() adapter.TaskConsumer {
|
|
return c.taskAdapter
|
|
}
|
|
|
|
func (c *compositeAdapter) EventProducer() adapter.EventProducer {
|
|
return c.eventAdapter
|
|
}
|
|
|
|
func (c *compositeAdapter) EventConsumer() adapter.EventConsumer {
|
|
return c.eventAdapter
|
|
}
|
|
|
|
func (c *compositeAdapter) Close() error {
|
|
var firstErr error
|
|
if err := c.eventAdapter.Close(); err != nil && firstErr == nil {
|
|
firstErr = err
|
|
}
|
|
if err := c.taskAdapter.Close(); err != nil && firstErr == nil {
|
|
firstErr = err
|
|
}
|
|
return firstErr
|
|
}
|