-
Notifications
You must be signed in to change notification settings - Fork 33
/
kafka_consumer.go
74 lines (62 loc) · 1.22 KB
/
kafka_consumer.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
package kafka
import (
"context"
"sync"
mq_ "github.com/kaydxh/golang/pkg/mq"
kafka "github.com/segmentio/kafka-go"
"github.com/sirupsen/logrus"
)
type Consumer struct {
*kafka.Reader
config kafka.ReaderConfig
msgCh chan mq_.Message
streamOnce sync.Once
closeOnce sync.Once
closeCh chan struct{}
}
func NewConsumer(config kafka.ReaderConfig) (*Consumer, error) {
c := &Consumer{
config: config,
msgCh: make(chan mq_.Message, 1024),
closeCh: make(chan struct{}),
}
r := kafka.NewReader(config)
c.Reader = r
return c, nil
}
func (c *Consumer) Topic() string {
return c.config.Topic
}
func (c *Consumer) ReadStream(ctx context.Context) <-chan mq_.Message {
c.streamOnce.Do(func() {
go func() {
for {
select {
case <-c.closeCh:
err := c.Reader.Close()
if err != nil {
logrus.WithError(err).Errorf("failed to close consumer")
}
if c.msgCh != nil {
close(c.msgCh)
}
return
default:
msg, err := c.Reader.ReadMessage(ctx)
c.msgCh <- KafkaMessage{
Err: err,
Msg: &msg,
}
}
}
}()
})
return c.msgCh
}
func (c *Consumer) Close() {
c.closeOnce.Do(func() {
if c.closeCh != nil {
close(c.closeCh)
}
})
}