-
Notifications
You must be signed in to change notification settings - Fork 0
xpulsar
omeyang edited this page Sep 17, 2026
·
1 revision
稳定性:Stable · 覆盖率:99.5% · 源码:
pkg/mq/xpulsar
Pulsar 客户端。基于 pulsar-client-go,加 DLQ + OTel 链路追踪。
- 多租户消息系统
- 延迟消息 / 定时消息
- 多订阅模式(Exclusive/Shared/Failover/KeyShared)
import "github.com/omeyang/xkit/pkg/mq/xpulsar"
c, _ := xpulsar.NewClient(xpulsar.Config{
URL: "pulsar://localhost:6650",
})
defer c.Close()
p, _ := c.NewProducer("topic-1")
_ = p.Send(ctx, []byte(payload))
con, _ := c.NewConsumer(xpulsar.ConsumerConfig{
Topic: "topic-1",
Subscription: "my-sub",
Type: xpulsar.SubscriptionShared,
})
_ = con.ConsumeLoop(ctx, handler)| 名称 | 说明 |
|---|---|
Client / Producer / Consumer |
三大类型 |
ConsumerConfig.Type |
Exclusive / Shared / Failover / KeyShared |
DLQConfig |
死信队列 |
TracingClient |
OTel 装饰器 |
- 覆盖率最高(99.5%):测试用例完整
- OTel propagator nil 守卫:与 xkafka 同样(OTelTracer 零值)
- 模式:Fail-Fast 消费循环
- 内部:mqcore
- API:api.md#pkgmqxpulsar