-
Notifications
You must be signed in to change notification settings - Fork 0
xkafka
omeyang edited this page Sep 17, 2026
·
1 revision
稳定性:Stable · 覆盖率:88.5% · 源码:
pkg/mq/xkafka
Kafka 客户端。基于 confluent-kafka-go,加 DLQ + OTel 链路追踪。
- 生产 / 消费 Kafka 消息
- DLQ 自动投递无法处理的消息
- 链路追踪自动传播
import "github.com/omeyang/xkit/pkg/mq/xkafka"
// 生产
p, _ := xkafka.NewProducer(xkafka.ProducerConfig{
Brokers: []string{"localhost:9092"},
})
defer p.Close()
_ = p.Produce(ctx, "topic-1", []byte(payload))
// 消费(fail-fast:handler panic 不补 recover)
c, _ := xkafka.NewConsumer(xkafka.ConsumerConfig{
Brokers: []string{"localhost:9092"},
Topic: "topic-1",
Group: "my-group",
DLQ: xkafka.DLQConfig{Topic: "topic-1-dlq", MaxRetries: 3},
})
_ = c.ConsumeLoop(ctx, func(ctx context.Context, msg *xkafka.Message) error {
return process(msg)
})| 名称 | 说明 |
|---|---|
Producer / NewProducer(cfg) |
生产 |
Consumer / NewConsumer(cfg) |
消费 |
ConsumeLoop(ctx, handler) error |
阻塞消费(fail-fast,handler panic 暴露不补 recover) |
Message |
消息封装 |
TracingProducer / TracingConsumer |
OTel 装饰器 |
DLQConfig |
DLQ topic + MaxRetries |
-
ctx取消前置检查(FG-M fix):sendToDLQInternal/redeliverMessage在 Produce 前检查ctx.Err(),防止 ctx 已取消时仍入队但不提交 offset -
setHeader重复 Header 修复(FG-M fix):删除所有同名 header 后追加新值,避免重复 traceparent -
TracingProducer/Consumer生命周期竞态(对抗审查 fix):避免 close 后还在 publish -
Fail-Fast:
ConsumeLoophandler panic 不补 recover(共享设计internal/mqcore)
- 模式:Fail-Fast 消费循环
- 内部:mqcore(共享底层)
- API:api.md#pkgmqxkafka