-
Notifications
You must be signed in to change notification settings - Fork 50
/
consumer.go
159 lines (120 loc) · 4.14 KB
/
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
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
package stream
import (
"context"
"fmt"
"sync/atomic"
"github.com/justtrackio/gosoline/pkg/cfg"
"github.com/justtrackio/gosoline/pkg/kernel"
"github.com/justtrackio/gosoline/pkg/log"
)
type ConsumerCallbackFactory func(ctx context.Context, config cfg.Config, logger log.Logger) (ConsumerCallback, error)
//go:generate mockery --name ConsumerCallback
type ConsumerCallback interface {
BaseConsumerCallback
Consume(ctx context.Context, model interface{}, attributes map[string]string) (bool, error)
}
//go:generate mockery --name RunnableConsumerCallback
type RunnableConsumerCallback interface {
ConsumerCallback
RunnableCallback
}
type Consumer struct {
*baseConsumer
callback ConsumerCallback
}
func NewConsumer(name string, callbackFactory ConsumerCallbackFactory) func(ctx context.Context, config cfg.Config, logger log.Logger) (kernel.Module, error) {
return func(ctx context.Context, config cfg.Config, logger log.Logger) (kernel.Module, error) {
loggerCallback := logger.WithChannel("consumerCallback")
contextEnforcingLogger := log.NewContextEnforcingLogger(loggerCallback)
var err error
var callback ConsumerCallback
var baseConsumer *baseConsumer
if callback, err = callbackFactory(ctx, config, contextEnforcingLogger); err != nil {
return nil, fmt.Errorf("can not initiate callback for consumer %s: %w", name, err)
}
contextEnforcingLogger.Enable()
if baseConsumer, err = NewBaseConsumer(ctx, config, logger, name, callback); err != nil {
return nil, fmt.Errorf("can not initiate base consumer: %w", err)
}
return NewConsumerWithInterfaces(baseConsumer, callback), nil
}
}
func NewConsumerWithInterfaces(base *baseConsumer, callback ConsumerCallback) *Consumer {
consumer := &Consumer{
baseConsumer: base,
callback: callback,
}
return consumer
}
func (c *Consumer) Run(kernelCtx context.Context) error {
return c.baseConsumer.run(kernelCtx, c.readData)
}
func (c *Consumer) readData(ctx context.Context) error {
defer c.logger.Debug("read from input is ending")
defer c.wg.Done()
for {
select {
case <-ctx.Done():
return nil
case cdata, ok := <-c.data:
if !ok {
return nil
}
if _, ok := cdata.msg.Attributes[AttributeAggregate]; ok {
c.processAggregateMessage(ctx, cdata)
} else {
c.processSingleMessage(ctx, cdata)
}
}
}
}
func (c *Consumer) processAggregateMessage(ctx context.Context, cdata *consumerData) {
var err error
start := c.clock.Now()
batch := make([]*Message, 0)
if ctx, _, err = c.encoder.Decode(ctx, cdata.msg, &batch); err != nil {
c.handleError(ctx, err, "an error occurred during disaggregation of the message")
return
}
c.Acknowledge(ctx, cdata, true)
for _, m := range batch {
_ = c.process(ctx, m, false) // we can't natively retry aggregate messages
}
duration := c.clock.Now().Sub(start)
atomic.AddInt32(&c.processed, int32(len(batch)))
c.writeMetricDurationAndProcessedCount(duration, len(batch))
}
func (c *Consumer) processSingleMessage(ctx context.Context, cdata *consumerData) {
start := c.clock.Now()
ack := c.process(ctx, cdata.msg, c.hasNativeRetry())
c.Acknowledge(ctx, cdata, ack)
duration := c.clock.Now().Sub(start)
atomic.AddInt32(&c.processed, 1)
c.writeMetricDurationAndProcessedCount(duration, 1)
}
func (c *Consumer) process(ctx context.Context, msg *Message, hasNativeRetry bool) bool {
defer c.recover(ctx, msg)
var err error
var ack bool
var model interface{}
var attributes map[string]string
if model = c.callback.GetModel(msg.Attributes); model == nil {
err := fmt.Errorf("can not get model for message attributes %v", msg.Attributes)
c.handleError(ctx, err, "an error occurred during the consume operation")
return false
}
if ctx, attributes, err = c.encoder.Decode(ctx, msg, model); err != nil {
c.handleError(ctx, err, "an error occurred during the consume operation")
return false
}
ctx, span := c.tracer.StartSpanFromContext(ctx, c.id)
defer span.Finish()
ctx = log.InitContext(ctx)
if ack, err = c.callback.Consume(ctx, model, attributes); err != nil {
c.handleError(ctx, err, "an error occurred during the consume operation")
}
if !ack && !hasNativeRetry {
c.retry(ctx, msg)
}
return ack
}