/
rebelance_handler.go
95 lines (79 loc) · 2.44 KB
/
rebelance_handler.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
package kstream
import (
"context"
"fmt"
"github.com/tryfix/kstream/consumer"
"github.com/tryfix/log"
)
type reBalanceHandler struct {
userHandler consumer.ReBalanceHandler
processors *processorPool
logger log.Logger
builder *StreamBuilder
rebalancedCount int
}
func (s *reBalanceHandler) OnPartitionRevoked(ctx context.Context, revoked []consumer.TopicPartition) error {
s.logger.Info(fmt.Sprintf(`partitions %v revoking...`, revoked))
defer s.logger.Info(fmt.Sprintf(`partitions %v revoked`, revoked))
for _, tp := range revoked {
s.processors.Processor(tp).Stop()
}
if s.userHandler != nil {
return s.userHandler.OnPartitionRevoked(ctx, revoked)
}
if err := s.startChangelogReplicas(revoked); err != nil {
return err
}
return nil
}
func (s *reBalanceHandler) OnPartitionAssigned(ctx context.Context, assigned []consumer.TopicPartition) error {
s.logger.Info(fmt.Sprintf(`partitions %v assigning...`, assigned))
defer s.logger.Info(fmt.Sprintf(`partitions %v assigned`, assigned))
for _, tp := range assigned {
if err := s.processors.addProcessor(tp); err != nil {
return err
}
if err := s.processors.Processor(tp).boot(); err != nil {
return err
}
if err := s.stopChangelogReplicas(assigned); err != nil {
return err
}
}
s.logger.Info(`streams assigned`)
if s.userHandler != nil {
return s.userHandler.OnPartitionAssigned(ctx, assigned)
}
s.rebalancedCount++
return nil
}
func (s *reBalanceHandler) stopChangelogReplicas(allocated []consumer.TopicPartition) error {
if len(allocated) > 0 && s.rebalancedCount > 0 {
for _, tp := range allocated {
// stop started replicas
if s.builder.streams[tp.Topic].config.changelog.replicated {
if err := s.builder.changelogReplicaManager.StopReplicas([]consumer.TopicPartition{
{Topic: s.builder.streams[tp.Topic].config.changelog.topic.Name, Partition: tp.Partition},
}); err != nil {
return err
}
}
}
}
return nil
}
func (s *reBalanceHandler) startChangelogReplicas(allocated []consumer.TopicPartition) error {
if len(allocated) > 0 {
for _, tp := range allocated {
// stop started replicas
if s.builder.streams[tp.Topic].config.changelog.replicated {
if err := s.builder.changelogReplicaManager.StartReplicas([]consumer.TopicPartition{
{Topic: s.builder.streams[tp.Topic].config.changelog.topic.Name, Partition: tp.Partition},
}); err != nil {
return err
}
}
}
}
return nil
}