mq

package
v0.8.2 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 15, 2026 License: Apache-2.0 Imports: 6 Imported by: 0

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)、分区键(Key)、broker 选型都是 policy。 OTel trace 透传(Headers 承载 W3C TraceContext)见子包 pkg/mq/otelmq(opt-in,对齐 franz-go kotel 的 publish/process 语义,但与具体 broker 解耦)。

Index

Constants

This section is empty.

Variables

View Source
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。

func (*Consumer) Start

func (c *Consumer) Start(ctx context.Context) error

Start 订阅所有登记的 topic,然后阻塞到 ctx 取消——满足 beauty.Service。

func (*Consumer) String

func (c *Consumer) String() string

String 满足 beauty.Service。

type ConsumerOption

type ConsumerOption func(*Consumer)

ConsumerOption 配置 Consumer。

func WithConsumerName

func WithConsumerName(name string) ConsumerOption

WithConsumerName 设置消费者名(日志/标识用)。

type Handler

type Handler func(ctx context.Context, msg Message) error

Handler 处理一条消息。返回非 nil error 表示处理失败——具体后果(重投/丢弃)由 broker 与 中间件决定。

func Chain

func Chain(h Handler, mw ...HandlerMiddleware) Handler

Chain 按声明顺序把中间件套到 h 外层(第一个最外层,最先执行)。

type HandlerMiddleware

type HandlerMiddleware func(Handler) Handler

HandlerMiddleware 包装 Handler(如重试、恢复 panic、埋点)。

func Recover

func Recover() HandlerMiddleware

Recover 把 handler 里的 panic 转成 error,避免打崩投递 goroutine。

func Retry

func Retry(attempts int, delay time.Duration) HandlerMiddleware

Retry 对返回 error 的处理重试至多 attempts 次,第 i 次失败后等 delay*(i+1)(线性退避); ctx 取消则立即返回。适合兜**瞬时**错误(下游抖动),不是持久化重投的替代。

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 中间件兜瞬时错误)。

func NewInProc

func NewInProc(opts ...InProcOption) *InProc

NewInProc 创建进程内 broker。

func (*InProc) Close

func (b *InProc) Close() error

Close 关闭 broker:拒绝后续发布/订阅。已存在的订阅由各自 ctx 解除。幂等。

func (*InProc) Publish

func (b *InProc) Publish(ctx context.Context, msg Message) error

Publish 把消息投给 topic 的订阅者:非 group 订阅者各投一份(扇出),每个 group 轮询选一个投递。

func (*InProc) Subscribe

func (b *InProc) Subscribe(ctx context.Context, topic string, h Handler, opts ...SubscribeOption) error

Subscribe 注册订阅;ctx 取消即解除该订阅并停止其投递 goroutine。

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 WithBuffer

func WithBuffer(n int) SubscribeOption

WithBuffer 设置该订阅的投递缓冲大小。

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 实现。

Directories

Path Synopsis
Package otelmq 为 pkg/mq 提供 OpenTelemetry trace 透传(opt-in)。
Package otelmq 为 pkg/mq 提供 OpenTelemetry trace 透传(opt-in)。

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL