消息队列

mq 包提供发布/订阅(pub/sub)的统一抽象,参考 Dapr Building Block 设计。

概念说明
Message消息体(Topic + Payload + Headers)
Handler消息处理函数(返回 nil=ack,error=nack)
Publisher发布者接口(Publish + Close)
Subscriber订阅者接口(Subscribe + Close)
Broker完整代理(同时实现 Publisher + Subscriber)
memory.Broker内置实现:channel fan-out,无缓冲反压,自动 baggage 注入/提取
MQComponentcomponents 适配器:声明式注册 + 自动启停

设计权衡

维度选择理由
抽象层级只抽象 topic / payload / headers屏蔽 Kafka partition / RabbitMQ exchange 等厂商专属语义
ack 语义Handler 返回 error = nack不同实现可映射到不同动作(memory 走 ErrorHandler,Kafka 不 commit offset)
内置实现进程内 channel + 无缓冲零依赖、保证不丢消息,反压慢消费者(牺牲吞吐换可靠)
持久化不持久化用作单进程事件总线 / 测试 mock;生产用 plugins
Baggage 传播Publish 自动注入 msg.Headers[“baggage”],handler 自动 extract全链路 tenant.id / cluster 等 K-V 透传
并发模型每订阅者独立 goroutinefan-out 隔离故障

使用方式

 1import (
 2    "github.com/go-zeus/zeus/components"
 3    "github.com/go-zeus/zeus/mq"
 4    "github.com/go-zeus/zeus/mq/memory"
 5)
 6
 7// 1. 直接使用 Broker(无 components)
 8broker := memory.New()
 9defer broker.Close()
10
11_ = broker.Subscribe(ctx, "orders.created", func(ctx context.Context, msg *mq.Message) error {
12    return nil
13})
14_ = broker.Publish(ctx, "orders.created", &mq.Message{Payload: []byte("order-1")})
15
16// 2. 自动装配
17app := components.NewApp(
18    components.NewMQComponent(memory.New()),
19    components.NewMQSubscription("orders.created", handleOrder),
20    components.NewMQSubscription("log.all", handleLog),
21)
22app.Run()

Baggage 自动传播

位置行为
Publish 出口自动 InjectMetadata(ctx, msg.Headers):ctx baggage → msg.Headers["baggage"](W3C 编码)
handler 入口自动 ExtractMetadata(ctx, msg.Headers)msg.Headers["baggage"] → handler ctx
Handler 内读取propagation.Get(ctx, "tenant.id") 直接拿到

完整示例参见 examples/15-mq/:3 个订阅者(不同 topic)+ baggage 全链路传播 + 优雅关闭。