-
Notifications
You must be signed in to change notification settings - Fork 0
/
async_consumer.go
84 lines (73 loc) · 2.19 KB
/
async_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
75
76
77
78
79
80
81
82
83
84
package scratch
import (
"context"
"log"
"sync"
"github.com/Shopify/sarama"
)
type AsyncConsumer struct {
sync.Mutex
ConsumerGroup sarama.ConsumerGroup
HandleFn func(context.Context, []byte) error
MaxCount uint32
MsgCount uint32
MsgLenSum uint32
Closed bool
}
func (c *AsyncConsumer) Handle(ctx context.Context, data []byte) error {
if c.HandleFn == nil {
return nil
}
if err := c.HandleFn(context.Background(), data); err != nil {
return err
}
c.MsgCount++
c.MsgLenSum += uint32(len(data))
return nil
}
func (c *AsyncConsumer) IsDone() bool {
return c.MsgCount >= c.MaxCount
}
// Setup is run at the beginning of a new session, before ConsumeClaim
func (c *AsyncConsumer) Setup(sarama.ConsumerGroupSession) error {
log.Println("async consumer setup...")
return nil
}
// Cleanup is run at the end of a session, once all ConsumeClaim goroutines have exited
func (c *AsyncConsumer) Cleanup(sarama.ConsumerGroupSession) error {
log.Println("async consumer cleanup...")
return nil
}
// ConsumeClaim must start a consumer loop of ConsumerGroupClaim's Messages().
func (c *AsyncConsumer) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
// NOTE:
// Do not move the code below to a goroutine.
// The `ConsumeClaim` itself is called within a goroutine, see:
// https://github.com/Shopify/sarama/blob/main/consumer_group.go#L27-L29
for {
select {
case message := <-claim.Messages():
if c.IsDone() {
return nil
}
go func() {
// log.Printf("Message claimed: value = %s, timestamp = %v, topic = %s", string(message.Value), message.Timestamp, message.Topic)
if err := c.Handle(session.Context(), message.Value); err != nil {
log.Println("consumer handle error:", err)
}
session.MarkMessage(message, "")
if !c.Closed && c.IsDone() {
c.Lock()
c.ConsumerGroup.Close()
c.Closed = true
c.Unlock()
}
}()
// Should return when `session.Context()` is done.
// If not, will raise `ErrRebalanceInProgress` or `read tcp <ip>:<port>: i/o timeout` when kafka rebalance. see:
// https://github.com/Shopify/sarama/issues/1192
case <-session.Context().Done():
return nil
}
}
}