-
Notifications
You must be signed in to change notification settings - Fork 0
/
subscriber.go
45 lines (37 loc) · 1.14 KB
/
subscriber.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
package pubsubx
import (
"context"
"time"
"github.com/ThreeDotsLabs/watermill/message"
)
type Subscriber interface {
// Subscribe subscribes to the topic.
Subscribe(ctx context.Context, topic string) (<-chan *message.Message, error)
// Close closes the subscriber.
Close() error
}
type BatchConsumerOptions struct {
// MaxBatchSize max amount of elements the batch will contain.
// Default value is 100 if nothing is specified.
MaxBatchSize int16
// MaxWaitTime max time that it will be waited until MaxBatchSize elements are received.
// Default value is 100ms if nothing is specified.
MaxWaitTime time.Duration
}
type subscriberOptions struct {
consumerModel ConsumerModel
batchConsumerOptions *BatchConsumerOptions
}
type SubscriberOption func(*subscriberOptions)
func WithDefaultConsumerModel() SubscriberOption {
return func(o *subscriberOptions) {
o.consumerModel = ConsumerModelDefault
o.batchConsumerOptions = nil
}
}
func WithBatchConsumerModel(batchOptions *BatchConsumerOptions) SubscriberOption {
return func(o *subscriberOptions) {
o.consumerModel = ConsumerModelBatch
o.batchConsumerOptions = batchOptions
}
}