diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionRuntimeManager.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionRuntimeManager.java index 3eba9dabe2427..48c63af7829b4 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionRuntimeManager.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionRuntimeManager.java @@ -813,7 +813,6 @@ private void addAssignment(Assignment assignment) { } private void startFunctionInstance(Assignment assignment) { - log.info("infos: {}", functionRuntimeInfos.getAll()); String fullyQualifiedInstanceId = FunctionCommon.getFullyQualifiedInstanceId(assignment.getInstance()); FunctionRuntimeInfo functionRuntimeInfo = _getFunctionRuntimeInfo(fullyQualifiedInstanceId); diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/MembershipManager.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/MembershipManager.java index 749f3a3bf9e0c..2bb56133c53bb 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/MembershipManager.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/MembershipManager.java @@ -208,6 +208,7 @@ public void checkFailures(FunctionMetaDataManager functionMetaDataManager, // check unassigned Collection needSchedule = new LinkedList<>(); Collection needRemove = new LinkedList<>(); + Map numRemoved = new HashMap<>(); for (Map.Entry entry : this.unsignedFunctionDurations.entrySet()) { Function.Instance instance = entry.getKey(); long unassignedDurationMs = entry.getValue(); @@ -217,6 +218,12 @@ public void checkFailures(FunctionMetaDataManager functionMetaDataManager, Function.Assignment assignment = assignmentMap.get(FunctionCommon.getFullyQualifiedInstanceId(instance)); if (assignment != null) { needRemove.add(assignment); + + Integer count = numRemoved.get(assignment.getWorkerId()); + if (count == null) { + count = 0; + } + numRemoved.put(assignment.getWorkerId(), count + 1); } triggerScheduler = true; } @@ -225,7 +232,8 @@ public void checkFailures(FunctionMetaDataManager functionMetaDataManager, functionRuntimeManager.removeAssignments(needRemove); } if (triggerScheduler) { - log.info("Functions that need scheduling/rescheduling: {}", needSchedule); + log.info("Failure check - Total number of instances that need to be scheduled/rescheduled: {} | Number of unassigned instances that need to be scheduled: {} | Number of instances on dead workers that need to be reassigned {}", + needSchedule.size(), needSchedule.size() - needRemove.size(), numRemoved); schedulerManager.schedule(); } } diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/SchedulerManager.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/SchedulerManager.java index 374512acea01f..873a26b41e44b 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/SchedulerManager.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/SchedulerManager.java @@ -18,10 +18,14 @@ */ package org.apache.pulsar.functions.worker; +import com.fasterxml.jackson.core.JsonProcessingException; import com.google.common.annotations.VisibleForTesting; +import com.google.common.base.Preconditions; import com.google.common.collect.Lists; import com.google.common.util.concurrent.ThreadFactoryBuilder; import io.netty.util.concurrent.DefaultThreadFactory; +import lombok.Builder; +import lombok.Data; import lombok.Getter; import lombok.Setter; import lombok.extern.slf4j.Slf4j; @@ -34,6 +38,7 @@ import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; +import org.apache.pulsar.common.util.ObjectMapperFactory; import org.apache.pulsar.common.util.Reflections; import org.apache.pulsar.functions.proto.Function; import org.apache.pulsar.functions.proto.Function.Assignment; @@ -239,7 +244,8 @@ public Future rebalanceIfNotInprogress() { @VisibleForTesting void invokeScheduler() { - + long startTime = System.nanoTime(); + Set currentMembership = membershipManager.getCurrentMembership() .stream().map(workerInfo -> workerInfo.getWorkerId()).collect(Collectors.toSet()); @@ -248,11 +254,13 @@ void invokeScheduler() { Map> workerIdToAssignments = functionRuntimeManager .getCurrentAssignments(); + // initialize stats collection + SchedulerStats schedulerStats = new SchedulerStats(workerIdToAssignments, currentMembership); + //delete assignments of functions and instances that don't exist anymore Iterator>> it = workerIdToAssignments.entrySet().iterator(); while (it.hasNext()) { Map.Entry> workerIdToAssignmentEntry = it.next(); - String workerId = workerIdToAssignmentEntry.getKey(); Map functionMap = workerIdToAssignmentEntry.getValue(); // remove instances that don't exist anymore @@ -268,6 +276,8 @@ void invokeScheduler() { functionRuntimeManager.deleteAssignment(fullyQualifiedInstanceId); // update message id associated with current view of assignments map lastMessageProduced = messageId; + // update stats + schedulerStats.removedAssignment(assignment); } return deleted; }); @@ -288,6 +298,8 @@ void invokeScheduler() { functionRuntimeManager.processAssignment(newAssignment); // update message id associated with current view of assignments map lastMessageProduced = messageId; + //update stats + schedulerStats.updatedAssignment(newAssignment); } if (functionMap.isEmpty()) { it.remove(); @@ -331,16 +343,26 @@ void invokeScheduler() { functionRuntimeManager.processAssignment(assignment); // update message id associated with current view of assignments map lastMessageProduced = messageId; + // update stats + schedulerStats.newAssignment(assignment); } + + log.info("Schedule summary - execution time: {} sec | total unassigned: {} | stats: {}\n{}", + (System.nanoTime() - startTime) / Math.pow(10, 9), + unassignedInstances.getLeft().size(), schedulerStats.getSummary(), schedulerStats); } private void invokeRebalance() { + long startTime = System.nanoTime(); Set currentMembership = membershipManager.getCurrentMembership() .stream().map(workerInfo -> workerInfo.getWorkerId()).collect(Collectors.toSet()); Map> workerIdToAssignments = functionRuntimeManager.getCurrentAssignments(); + // initialize stats collection + SchedulerStats schedulerStats = new SchedulerStats(workerIdToAssignments, currentMembership); + // filter out assignments of workers that are not currently in the active membership List currentAssignments = workerIdToAssignments .entrySet() @@ -368,8 +390,13 @@ private void invokeRebalance() { functionRuntimeManager.processAssignment(assignment); // update message id associated with current view of assignments map lastMessageProduced = messageId; + // update stats + schedulerStats.newAssignment(assignment); } - log.info("Rebalance - Total number of new assignments computed: {}", rebalancedAssignments.size()); + + log.info("Rebalance summary - execution time: {} sec | stats: {}\n{}", + (System.nanoTime() - startTime) / Math.pow(10, 9), schedulerStats.getSummary(), schedulerStats); + rebalanceInProgess.set(false); } @@ -526,4 +553,117 @@ static String checkHeartBeatFunction(Instance funInstance) { public static class RebalanceInProgressException extends RuntimeException { } + + private static class SchedulerStats { + + @Builder + @Data + private static class WorkerStats { + private int originalNumAssignments; + private int finalNumAssignments; + private int instancesAdded; + private int instancesRemoved; + private int instancesUpdated; + private boolean alive; + } + + private Map workerStatsMap = new HashMap<>(); + + private Map instanceToWorkerId = new HashMap<>(); + + public SchedulerStats(Map> workerIdToAssignments, Set workers) { + + for(String workerId : workers) { + WorkerStats.WorkerStatsBuilder workerStats = WorkerStats.builder().alive(true); + Map assignmentMap = workerIdToAssignments.get(workerId); + if (assignmentMap != null) { + workerStats.originalNumAssignments(assignmentMap.size()); + workerStats.finalNumAssignments(assignmentMap.size()); + + for (String fullyQualifiedInstanceId : assignmentMap.keySet()) { + instanceToWorkerId.put(fullyQualifiedInstanceId, workerId); + } + } else { + workerStats.originalNumAssignments(0); + workerStats.finalNumAssignments(0); + } + + workerStatsMap.put(workerId, workerStats.build()); + } + + // workers with assignments that are dead + for (Map.Entry> entry : workerIdToAssignments.entrySet()) { + String workerId = entry.getKey(); + Map assignmentMap = entry.getValue(); + if (!workers.contains(workerId)) { + WorkerStats workerStats = WorkerStats.builder() + .alive(false) + .originalNumAssignments(assignmentMap.size()) + .finalNumAssignments(assignmentMap.size()) + .build(); + workerStatsMap.put(workerId, workerStats); + } + } + } + + public void removedAssignment(Assignment assignment) { + String workerId = assignment.getWorkerId(); + WorkerStats stats = workerStatsMap.get(workerId); + Preconditions.checkNotNull(stats); + + stats.instancesRemoved++; + stats.finalNumAssignments--; + } + + public void newAssignment(Assignment assignment) { + String fullyQualifiedInstanceId = FunctionCommon.getFullyQualifiedInstanceId(assignment.getInstance()); + String newWorkerId = assignment.getWorkerId(); + String oldWorkerId = instanceToWorkerId.get(fullyQualifiedInstanceId); + if (oldWorkerId != null) { + WorkerStats oldWorkerStats = workerStatsMap.get(oldWorkerId); + Preconditions.checkNotNull(oldWorkerStats); + + oldWorkerStats.instancesRemoved++; + oldWorkerStats.finalNumAssignments--; + } + + WorkerStats newWorkerStats = workerStatsMap.get(newWorkerId); + Preconditions.checkNotNull(newWorkerStats); + + newWorkerStats.instancesAdded++; + newWorkerStats.finalNumAssignments++; + } + + public void updatedAssignment(Assignment assignment) { + String workerId = assignment.getWorkerId(); + WorkerStats stats = workerStatsMap.get(workerId); + Preconditions.checkNotNull(stats); + + stats.instancesUpdated++; + } + + public String getSummary() { + int totalAdded = 0; + int totalUpdated = 0; + int totalRemoved = 0; + + for (Map.Entry entry : workerStatsMap.entrySet()) { + WorkerStats workerStats = entry.getValue(); + totalAdded += workerStats.instancesAdded; + totalUpdated += workerStats.instancesUpdated; + totalRemoved += workerStats.instancesRemoved; + } + + return String.format("{\"Added\": %d, \"Updated\": %d, \"removed\": %d}", totalAdded, totalUpdated, totalRemoved); + } + + @Override + public String toString() { + try { + return ObjectMapperFactory.getThreadLocal().writerWithDefaultPrettyPrinter().writeValueAsString(workerStatsMap); + } catch (JsonProcessingException e) { + throw new RuntimeException(e); + } + } + } } diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/scheduler/RoundRobinScheduler.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/scheduler/RoundRobinScheduler.java index 59b5137ed2ec4..0e8250643811e 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/scheduler/RoundRobinScheduler.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/scheduler/RoundRobinScheduler.java @@ -18,11 +18,7 @@ */ package org.apache.pulsar.functions.worker.scheduler; -import com.fasterxml.jackson.core.JsonProcessingException; -import lombok.Builder; -import lombok.Data; import lombok.extern.slf4j.Slf4j; -import org.apache.pulsar.common.util.ObjectMapperFactory; import org.apache.pulsar.functions.proto.Function.Assignment; import org.apache.pulsar.functions.proto.Function.Instance; @@ -41,9 +37,9 @@ public class RoundRobinScheduler implements IScheduler { @Override public List schedule(List unassignedFunctionInstances, - List currentAssignments, Set workers) { + List currentAssignments, Set workers) { - Map> workerIdToAssignment = new HashMap<>(); + Map> workerIdToAssignment = new HashMap<>(); List newAssignments = Lists.newArrayList(); for (String workerId : workers) { @@ -51,26 +47,26 @@ public List schedule(List unassignedFunctionInstances, } for (Assignment existingAssignment : currentAssignments) { - workerIdToAssignment.get(existingAssignment.getWorkerId()).add(existingAssignment); + workerIdToAssignment.get(existingAssignment.getWorkerId()).add(existingAssignment.getInstance()); } for (Instance unassignedFunctionInstance : unassignedFunctionInstances) { String workerId = findNextWorker(workerIdToAssignment); Assignment newAssignment = Assignment.newBuilder().setInstance(unassignedFunctionInstance) .setWorkerId(workerId).build(); - workerIdToAssignment.get(workerId).add(newAssignment); + workerIdToAssignment.get(workerId).add(newAssignment.getInstance()); newAssignments.add(newAssignment); } return newAssignments; } - private String findNextWorker(Map> workerIdToAssignment) { + private String findNextWorker(Map> workerIdToAssignment) { String targetWorkerId = null; int least = Integer.MAX_VALUE; - for (Map.Entry> entry : workerIdToAssignment.entrySet()) { + for (Map.Entry> entry : workerIdToAssignment.entrySet()) { String workerId = entry.getKey(); - List workerAssignments = entry.getValue(); + List workerAssignments = entry.getValue(); if (workerAssignments.size() < least) { targetWorkerId = workerId; least = workerAssignments.size(); @@ -82,7 +78,7 @@ private String findNextWorker(Map> workerIdToAssignment @Override public List rebalance(List currentAssignments, Set workers) { - Map> workerToAssignmentMap = new HashMap<>(); + Map> workerToAssignmentMap = new HashMap<>(); workers.forEach(workerId -> workerToAssignmentMap.put(workerId, new LinkedList<>())); @@ -90,14 +86,13 @@ public List rebalance(List currentAssignments, Set newAssignments = new LinkedList<>(); - RebalanceStats rebalanceStats = new RebalanceStats(workerToAssignmentMap); int iterations = 0; while(true) { iterations++; - Map.Entry> mostAssignmentsWorker = findWorkerWithMostAssignments(workerToAssignmentMap); + Map.Entry> mostAssignmentsWorker = findWorkerWithMostAssignments(workerToAssignmentMap); - Map.Entry> leastAssignmentsWorker = findWorkerWithLeastAssignments(workerToAssignmentMap); + Map.Entry> leastAssignmentsWorker = findWorkerWithLeastAssignments(workerToAssignmentMap); if (mostAssignmentsWorker.getValue().size() == leastAssignmentsWorker.getValue().size() || mostAssignmentsWorker.getValue().size() == leastAssignmentsWorker.getValue().size() + 1) { @@ -107,12 +102,8 @@ public List rebalance(List currentAssignments, Set src = workerToAssignmentMap.get(mostAssignmentsWorkerId); - Queue dest = workerToAssignmentMap.get(leastAssignmentsWorkerId); - - // update stats - rebalanceStats.decrementInstance(mostAssignmentsWorkerId); - rebalanceStats.incrementInstance(leastAssignmentsWorkerId); + Queue src = (Queue) workerToAssignmentMap.get(mostAssignmentsWorkerId); + Queue dest = (Queue) workerToAssignmentMap.get(leastAssignmentsWorkerId); Instance instance = src.poll(); Assignment newAssignment = Assignment.newBuilder() @@ -124,70 +115,18 @@ public List rebalance(List currentAssignments, Set> findWorkerWithLeastAssignments(Map> workerToAssignmentMap) { + private Map.Entry> findWorkerWithLeastAssignments(Map> workerToAssignmentMap) { return workerToAssignmentMap.entrySet().stream().min(Comparator.comparingInt(o -> o.getValue().size())).get(); } - private Map.Entry> findWorkerWithMostAssignments(Map> workerToAssignmentMap) { + private Map.Entry> findWorkerWithMostAssignments(Map> workerToAssignmentMap) { return workerToAssignmentMap.entrySet().stream().max(Comparator.comparingInt(o -> o.getValue().size())).get(); } - @Data - private static class RebalanceStats { - @Override - public String toString() { - try { - return ObjectMapperFactory.getThreadLocal().writerWithDefaultPrettyPrinter().writeValueAsString(workerStatsMap); - } catch (JsonProcessingException e) { - throw new RuntimeException(e); - } - } - - @Builder - @Data - private static class WorkerStats { - private int originalNumAssignments; - private int numAssignmentsAfterRebalance; - private int instancesAdded; - private int instancesRemoved; - } - - private Map workerStatsMap = new HashMap<>(); - - public RebalanceStats(Map> workerToAssignmentMap) { - for(Map.Entry> entry : workerToAssignmentMap.entrySet()) { - WorkerStats workerStats = WorkerStats.builder() - .originalNumAssignments(entry.getValue().size()) - .numAssignmentsAfterRebalance(entry.getValue().size()) - .build(); - workerStatsMap.put(entry.getKey(), workerStats); - } - } - - private void decrementInstance(String workerId) { - WorkerStats stats = workerStatsMap.get(workerId); - if (stats == null) { - throw new RuntimeException("Rebalance stats for worker " + workerId + " shouldn't be null"); - } - - stats.instancesRemoved++; - stats.numAssignmentsAfterRebalance--; - } - - private void incrementInstance(String workerId) { - WorkerStats stats = workerStatsMap.get(workerId); - if (stats == null) { - throw new RuntimeException("Rebalance stats for worker " + workerId + " shouldn't be null"); - } - - stats.instancesAdded++; - stats.numAssignmentsAfterRebalance++; - } - } }