-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathsubscriber.go
89 lines (76 loc) · 1.89 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
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
package repli
import (
"context"
"fmt"
"regexp"
"strings"
"github.com/redis/go-redis/v9"
log "github.com/sirupsen/logrus"
)
type KeyspaceEvent struct {
Key string
Action string
}
type Subscriber struct {
eventQueueSize int
keyPattern string
skipPatterns []*regexp.Regexp
subscriber *redis.Client
pubsub *redis.PubSub
C chan *KeyspaceEvent
}
func NewSubscriber(config *CommonConfig, eventQueueSize int) *Subscriber {
subscriber := redis.NewClient(&redis.Options{
Addr: config.SourceEndpoint,
DB: config.RedisDatabase,
MaxRetries: config.MaxRetries,
PoolSize: 1,
})
return &Subscriber{
eventQueueSize: eventQueueSize,
keyPattern: fmt.Sprintf("__keyspace@%d__:%s", config.RedisDatabase, config.KeyspacePattern),
skipPatterns: CompileRegExpPatterns(config.SkipKeyPatterns),
subscriber: subscriber,
C: make(chan *KeyspaceEvent),
}
}
func (s *Subscriber) Close() {
close(s.C)
if s.pubsub != nil {
s.pubsub.Close()
}
s.subscriber.Close()
}
func (s *Subscriber) Run(l *log.Entry, metrics *Metrics) {
ctx := context.Background()
s.pubsub = s.subscriber.PSubscribe(ctx, s.keyPattern)
// Wait for PSUBSCRIBE confirmation
_, err := s.pubsub.Receive(ctx)
if err != nil {
l.Fatal(err)
}
keyspaceEventCh := s.pubsub.Channel(redis.WithChannelSize(s.eventQueueSize))
loop:
for event := range keyspaceEventCh {
splits := strings.SplitN(event.Channel, ":", 2)
if len(splits) != 2 {
l.WithFields(log.Fields{
"eventPayload": event.Payload,
}).Error("unknown keyspace event")
continue
}
key := splits[1]
for _, re := range s.skipPatterns {
if re.Match([]byte(key)) {
l.WithFields(log.Fields{
"key": key,
}).Debug("skip pattern matched")
continue loop
}
}
action := event.Payload
metrics.Modify(key, action)
s.C <- &KeyspaceEvent{key, action}
metrics.Received()
}
}