|
@ -27,13 +27,8 @@ func AllocateTopicPartitions(brokers cmap.ConcurrentMap[string, *BrokerStats], p |
|
|
assignments = append(assignments, assignment) |
|
|
assignments = append(assignments, assignment) |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
// pick the brokers
|
|
|
|
|
|
pickedBrokers := pickBrokers(brokers, partitionCount) |
|
|
|
|
|
|
|
|
EnsureAssignmentsToActiveBrokers(brokers, 1, assignments) |
|
|
|
|
|
|
|
|
// assign the partitions to brokers
|
|
|
|
|
|
for i, assignment := range assignments { |
|
|
|
|
|
assignment.LeaderBroker = pickedBrokers[i] |
|
|
|
|
|
} |
|
|
|
|
|
glog.V(0).Infof("allocate topic partitions %d: %v", len(assignments), assignments) |
|
|
glog.V(0).Infof("allocate topic partitions %d: %v", len(assignments), assignments) |
|
|
return |
|
|
return |
|
|
} |
|
|
} |
|
@ -91,6 +86,8 @@ func pickBrokersExcluded(brokers []string, count int, excludedLeadBroker string, |
|
|
|
|
|
|
|
|
// EnsureAssignmentsToActiveBrokers ensures the assignments are assigned to active brokers
|
|
|
// EnsureAssignmentsToActiveBrokers ensures the assignments are assigned to active brokers
|
|
|
func EnsureAssignmentsToActiveBrokers(activeBrokers cmap.ConcurrentMap[string, *BrokerStats], followerCount int, assignments []*mq_pb.BrokerPartitionAssignment) (hasChanges bool) { |
|
|
func EnsureAssignmentsToActiveBrokers(activeBrokers cmap.ConcurrentMap[string, *BrokerStats], followerCount int, assignments []*mq_pb.BrokerPartitionAssignment) (hasChanges bool) { |
|
|
|
|
|
glog.V(0).Infof("EnsureAssignmentsToActiveBrokers: activeBrokers: %v, followerCount: %d, assignments: %v", activeBrokers.Count(), followerCount, assignments) |
|
|
|
|
|
|
|
|
candidates := make([]string, 0, activeBrokers.Count()) |
|
|
candidates := make([]string, 0, activeBrokers.Count()) |
|
|
for brokerStatsItem := range activeBrokers.IterBuffered() { |
|
|
for brokerStatsItem := range activeBrokers.IterBuffered() { |
|
|
candidates = append(candidates, brokerStatsItem.Key) |
|
|
candidates = append(candidates, brokerStatsItem.Key) |
|
@ -122,9 +119,11 @@ func EnsureAssignmentsToActiveBrokers(activeBrokers cmap.ConcurrentMap[string, * |
|
|
pickedBrokers := pickBrokersExcluded(candidates, count, assignment.LeaderBroker, assignment.FollowerBrokers) |
|
|
pickedBrokers := pickBrokersExcluded(candidates, count, assignment.LeaderBroker, assignment.FollowerBrokers) |
|
|
i := 0 |
|
|
i := 0 |
|
|
if assignment.LeaderBroker == "" { |
|
|
if assignment.LeaderBroker == "" { |
|
|
assignment.LeaderBroker = pickedBrokers[i] |
|
|
|
|
|
i++ |
|
|
|
|
|
hasChanges = true |
|
|
|
|
|
|
|
|
if i < len(pickedBrokers) { |
|
|
|
|
|
assignment.LeaderBroker = pickedBrokers[i] |
|
|
|
|
|
i++ |
|
|
|
|
|
hasChanges = true |
|
|
|
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
hasEmptyFollowers := false |
|
|
hasEmptyFollowers := false |
|
@ -158,5 +157,6 @@ func EnsureAssignmentsToActiveBrokers(activeBrokers cmap.ConcurrentMap[string, * |
|
|
|
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
glog.V(0).Infof("EnsureAssignmentsToActiveBrokers: activeBrokers: %v, followerCount: %d, assignments: %v hasChanges: %v", activeBrokers.Count(), followerCount, assignments, hasChanges) |
|
|
return |
|
|
return |
|
|
} |
|
|
} |