Documentation
¶
Overview ¶
Package mq 提供传输无关的消息队列抽象:发布/订阅接口 + 一个「消费者即 beauty.Service」 的运行原语 + 处理中间件。补齐框架跨服务异步的空白——此前只有进程内 eventbus(泛型扇出) 与 webhook(HTTP 推),没有面向真 broker(NATS/Kafka/…)的统一抽象。
分层:
- Publisher / Subscriber:传输无关接口,由具体 broker 实现。本包自带零依赖的进程内实现 (NewInProc),用于单体/开发/测试;真 broker 作为 opt-in 子包实现同一接口(如未来的 pkg/infra/nats),不强引依赖。
- Consumer:把一组 (topic, handler) 订阅包成 beauty.Service(Start/String/Ready), 随 app 优雅停机。
- HandlerMiddleware:Recover(吞 panic)、Retry(瞬时错误重试)等,Chain 组合。
语义:
- 订阅按 ctx 绑定生命周期:Subscribe 传入的 ctx 取消即解除该订阅(不影响 broker 其它订阅)。
- Group(队列组):同一 (topic, group) 的多个订阅者**竞争消费**(每条消息只投给组内一个), 用于多副本水平扩展;不设 group 则**扇出**(每个订阅者都收到)。语义对齐 NATS queue group / Kafka consumer group。
- 投递保证由 broker 决定:进程内实现是 at-most-once(handler 出错不重投,用 Retry 中间件兜 瞬时错误);要持久化/重投/exactly-once 用支持的 broker(如 JetStream)。
边界(机制而非策略):序列化(Body 是 []byte)、trace 透传(用 Headers 承载,配 pkg/metadata)、 分区键(Key)、broker 选型都是 policy。
Index ¶
Constants ¶
This section is empty.
Variables ¶
var ErrClosed = errors.New("mq: broker closed")
ErrClosed 表示 broker 已关闭,不再接受发布/订阅。
Functions ¶
This section is empty.
Types ¶
type Consumer ¶
type Consumer struct {
// contains filtered or unexported fields
}
Consumer 把一组订阅包成 beauty.Service:Start 时全部 Subscribe,随后阻塞到 ctx 取消 (各订阅随之解除)。结构上满足 beauty.Service(Start/String)+ ReadyNotifier(Ready)。 零值不可用,用 NewConsumer 构造;用 Handle 链式登记订阅。
func NewConsumer ¶
func NewConsumer(sub Subscriber, opts ...ConsumerOption) *Consumer
NewConsumer 创建消费者。sub 是任意 Subscriber 实现(进程内或 broker)。
func (*Consumer) Handle ¶
func (c *Consumer) Handle(topic string, h Handler, opts ...SubscribeOption) *Consumer
Handle 登记一个订阅(链式);在 Start 时统一 Subscribe。
func (*Consumer) Ready ¶
func (c *Consumer) Ready() <-chan struct{}
Ready 在所有订阅完成后关闭——满足 beauty.ReadyNotifier。
type ConsumerOption ¶
type ConsumerOption func(*Consumer)
ConsumerOption 配置 Consumer。
func WithConsumerName ¶
func WithConsumerName(name string) ConsumerOption
WithConsumerName 设置消费者名(日志/标识用)。
type Handler ¶
Handler 处理一条消息。返回非 nil error 表示处理失败——具体后果(重投/丢弃)由 broker 与 中间件决定。
func Chain ¶
func Chain(h Handler, mw ...HandlerMiddleware) Handler
Chain 按声明顺序把中间件套到 h 外层(第一个最外层,最先执行)。
type HandlerMiddleware ¶
HandlerMiddleware 包装 Handler(如重试、恢复 panic、埋点)。
func Recover ¶
func Recover() HandlerMiddleware
Recover 把 handler 里的 panic 转成 error,避免打崩投递 goroutine。
type InProc ¶
type InProc struct {
// contains filtered or unexported fields
}
InProc 是零依赖的进程内 broker,同时实现 Publisher 与 Subscriber——用于单体部署、开发 与测试。语义对齐真 broker:同 topic 下不设 group 的订阅者**扇出**(都收到),同一 (topic, group) 的订阅者**竞争消费**(轮询,每条只投一个)。投递经每订阅一条缓冲 channel + 一个投递 goroutine,不阻塞发布(缓冲满则背压)。at-most-once:handler 出错只记日志、不重投 (用 Retry 中间件兜瞬时错误)。
type InProcOption ¶
type InProcOption func(*InProc)
InProcOption 配置 InProc。
func WithDefaultBuffer ¶
func WithDefaultBuffer(n int) InProcOption
WithDefaultBuffer 设置订阅默认投递缓冲(未用 WithBuffer 指定时;默认 64)。
type Message ¶
type Message struct {
Topic string // 主题(broker 的 subject/topic)
Key string // 分区/有序键(可空;broker 用它做分区或有序投递)
Body []byte // 负载
Headers map[string]string // 元数据(content-type、trace 上下文等)
}
Message 是一条消息。Body 是已序列化的负载(编解码是 policy)。
type Publisher ¶
type Publisher interface {
// Publish 发布一条消息;应并发安全。
Publish(ctx context.Context, msg Message) error
}
Publisher 是发布侧,由 broker 实现。
type SubConfig ¶
type SubConfig struct {
Group string // 队列组:同组竞争消费;空则扇出
Buffer int // 投递缓冲(broker 实现可参考;<=0 用实现默认)
}
SubConfig 是订阅参数(供 broker 实现读取)。
func ApplySubOptions ¶
func ApplySubOptions(opts ...SubscribeOption) SubConfig
ApplySubOptions 供 broker 实现用:把 options 归一成 SubConfig。
type SubscribeOption ¶
type SubscribeOption func(*SubConfig)
SubscribeOption 配置一次订阅。
func WithGroup ¶
func WithGroup(group string) SubscribeOption
WithGroup 设置队列组(同 (topic, group) 竞争消费,用于多副本水平扩展)。
type Subscriber ¶
type Subscriber interface {
// Subscribe 为 topic 注册 handler,订阅生命周期与 ctx 绑定(ctx 取消即解除该订阅)。
// 通过 WithGroup 指定队列组做竞争消费。方法返回后投递即开始(在 broker 的 goroutine 上)。
Subscribe(ctx context.Context, topic string, h Handler, opts ...SubscribeOption) error
}
Subscriber 是订阅侧,由 broker 实现。