This repository has been archived by the owner on Oct 9, 2023. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 63
/
Copy pathrandom_cluster_selector.go
173 lines (159 loc) · 5.62 KB
/
random_cluster_selector.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
160
161
162
163
164
165
166
167
168
169
170
171
172
173
package impl
import (
"context"
"fmt"
"hash/fnv"
"math/rand"
"github.com/lyft/flyteadmin/pkg/executioncluster"
"github.com/lyft/flyteadmin/pkg/executioncluster/interfaces"
"github.com/lyft/flytestdlib/random"
runtime "github.com/lyft/flyteadmin/pkg/runtime/interfaces"
"github.com/lyft/flytestdlib/promutils"
)
// Implementation of Random cluster selector
// Selects cluster based on weights and domains.
type RandomClusterSelector struct {
domainWeightedRandomMap map[string]random.WeightedRandomList
executionTargetMap map[string]executioncluster.ExecutionTarget
}
func getRandSource(seed string) (rand.Source, error) {
h := fnv.New64a()
_, err := h.Write([]byte(seed))
if err != nil {
return nil, err
}
hashedSeed := int64(h.Sum64())
return rand.NewSource(hashedSeed), nil
}
func getValidDomainMap(validDomains runtime.DomainsConfig) map[string]runtime.Domain {
domainMap := make(map[string]runtime.Domain)
for _, domain := range validDomains {
domainMap[domain.ID] = domain
}
return domainMap
}
func getExecutionTargetMap(scope promutils.Scope, executionTargetProvider interfaces.ExecutionTargetProvider, clusterConfig runtime.ClusterConfiguration) (map[string]executioncluster.ExecutionTarget, error) {
executionTargetMap := make(map[string]executioncluster.ExecutionTarget)
for _, cluster := range clusterConfig.GetClusterConfigs() {
if _, ok := executionTargetMap[cluster.Name]; ok {
return nil, fmt.Errorf("duplicate clusters for name %s", cluster.Name)
}
executionTarget, err := executionTargetProvider.GetExecutionTarget(scope, cluster)
if err != nil {
return nil, err
}
executionTargetMap[cluster.Name] = *executionTarget
}
return executionTargetMap, nil
}
func getDomainsForCluster(cluster runtime.ClusterConfig, domainMap map[string]runtime.Domain) ([]string, error) {
if len(cluster.AllowedDomains) == 0 {
allDomains := make([]string, len(domainMap))
index := 0
for id := range domainMap {
allDomains[index] = id
index++
}
return allDomains, nil
}
for _, allowedDomain := range cluster.AllowedDomains {
if _, ok := domainMap[allowedDomain]; !ok {
return nil, fmt.Errorf("invalid domain %s", allowedDomain)
}
}
return cluster.AllowedDomains, nil
}
func getDomainWeightedRandomForCluster(ctx context.Context, scope promutils.Scope, executionTargetProvider interfaces.ExecutionTargetProvider,
clusterConfig runtime.ClusterConfiguration,
domainMap map[string]runtime.Domain) (map[string]random.WeightedRandomList, error) {
domainEntriesMap := make(map[string][]random.Entry)
for _, cluster := range clusterConfig.GetClusterConfigs() {
// If cluster is not enabled, it is not eligible for selection
if !cluster.Enabled {
continue
}
executionTarget, err := executionTargetProvider.GetExecutionTarget(scope, cluster)
if err != nil {
return nil, err
}
targetEntry := random.Entry{
Item: *executionTarget,
Weight: cluster.Weight,
}
clusterDomains, err := getDomainsForCluster(cluster, domainMap)
if err != nil {
return nil, err
}
for _, domain := range clusterDomains {
if _, ok := domainEntriesMap[domain]; ok {
domainEntriesMap[domain] = append(domainEntriesMap[domain], targetEntry)
} else {
domainEntriesMap[domain] = []random.Entry{targetEntry}
}
}
}
domainWeightedRandomMap := make(map[string]random.WeightedRandomList)
for domain, entries := range domainEntriesMap {
weightedRandomList, err := random.NewWeightedRandom(ctx, entries)
if err != nil {
return nil, err
}
domainWeightedRandomMap[domain] = weightedRandomList
}
return domainWeightedRandomMap, nil
}
func (s RandomClusterSelector) GetAllValidTargets() []executioncluster.ExecutionTarget {
v := make([]executioncluster.ExecutionTarget, 0)
for _, value := range s.executionTargetMap {
if value.Enabled {
v = append(v, value)
}
}
return v
}
func (s RandomClusterSelector) GetTarget(spec *executioncluster.ExecutionTargetSpec) (*executioncluster.ExecutionTarget, error) {
if spec == nil || (spec.TargetID == "" && spec.ExecutionID == nil) {
return nil, fmt.Errorf("invalid executionTargetSpec %v", spec)
}
if spec.TargetID != "" {
if val, ok := s.executionTargetMap[spec.TargetID]; ok {
return &val, nil
}
return nil, fmt.Errorf("invalid cluster target %s", spec.TargetID)
}
if spec.ExecutionID != nil {
if weightedRandomList, ok := s.domainWeightedRandomMap[spec.ExecutionID.GetDomain()]; ok {
executionName := spec.ExecutionID.GetName()
if executionName != "" {
randSrc, err := getRandSource(executionName)
if err != nil {
return nil, err
}
result, err := weightedRandomList.GetWithSeed(randSrc)
if err != nil {
return nil, err
}
execTarget := result.(executioncluster.ExecutionTarget)
return &execTarget, nil
}
execTarget := weightedRandomList.Get().(executioncluster.ExecutionTarget)
return &execTarget, nil
}
}
return nil, fmt.Errorf("invalid executionTargetSpec %v", *spec)
}
func NewRandomClusterSelector(scope promutils.Scope, clusterConfig runtime.ClusterConfiguration, executionTargetProvider interfaces.ExecutionTargetProvider, domainConfig *runtime.DomainsConfig) (interfaces.ClusterInterface, error) {
executionTargetMap, err := getExecutionTargetMap(scope, executionTargetProvider, clusterConfig)
if err != nil {
return nil, err
}
domainMap := getValidDomainMap(*domainConfig)
domainWeightedRandomMap, err := getDomainWeightedRandomForCluster(context.Background(), scope, executionTargetProvider, clusterConfig, domainMap)
if err != nil {
return nil, err
}
return &RandomClusterSelector{
domainWeightedRandomMap: domainWeightedRandomMap,
executionTargetMap: executionTargetMap,
}, nil
}