JIRA-1212 closes #94
diff --git a/giraph-core/src/main/java/org/apache/giraph/conf/GiraphConfiguration.java b/giraph-core/src/main/java/org/apache/giraph/conf/GiraphConfiguration.java index 17c48a5..832e8b6 100644 --- a/giraph-core/src/main/java/org/apache/giraph/conf/GiraphConfiguration.java +++ b/giraph-core/src/main/java/org/apache/giraph/conf/GiraphConfiguration.java
@@ -671,6 +671,16 @@ } /** + * How many mappers is job asking for, taking into account whether master + * is running on the same mapper as worker or not + * + * @return How many mappers is job asking for + */ + public final int getMaxMappers() { + return getMaxWorkers() + (SPLIT_MASTER_WORKER.get(this) ? 1 : 0); + } + + /** * Utilize an existing ZooKeeper service. If this is not set, ZooKeeper * will be dynamically started by Giraph for this job. *
diff --git a/giraph-core/src/main/java/org/apache/giraph/job/DefaultJobProgressTrackerService.java b/giraph-core/src/main/java/org/apache/giraph/job/DefaultJobProgressTrackerService.java index 243ab81..a1e7f12 100644 --- a/giraph-core/src/main/java/org/apache/giraph/job/DefaultJobProgressTrackerService.java +++ b/giraph-core/src/main/java/org/apache/giraph/job/DefaultJobProgressTrackerService.java
@@ -93,7 +93,7 @@ MAX_ALLOWED_TIME_WITHOUT_PROGRESS_MS.get(conf); CombinedWorkerProgress lastProgress = null; while (!finished) { - if (mappersStarted == conf.getMaxWorkers() + 1 && + if (mappersStarted == conf.getMaxMappers() && !workerProgresses.isEmpty()) { // Combine and log CombinedWorkerProgress combinedWorkerProgress = @@ -202,7 +202,7 @@ public synchronized void mapperStarted() { mappersStarted++; if (LOG.isInfoEnabled()) { - if (mappersStarted == conf.getMaxWorkers() + 1) { + if (mappersStarted == conf.getMaxMappers()) { LOG.info("Got all " + mappersStarted + " mappers"); jobGotAllMappers(); } else { @@ -210,7 +210,7 @@ UPDATE_MILLISECONDS) { lastTimeMappersStartedLogged = System.currentTimeMillis(); LOG.info("Got " + mappersStarted + " but needs " + - (conf.getMaxWorkers() + 1) + " mappers"); + conf.getMaxMappers() + " mappers"); } } }
diff --git a/giraph-core/src/main/java/org/apache/giraph/master/MasterThread.java b/giraph-core/src/main/java/org/apache/giraph/master/MasterThread.java index 8e4e0b8..0f6a3a0 100644 --- a/giraph-core/src/main/java/org/apache/giraph/master/MasterThread.java +++ b/giraph-core/src/main/java/org/apache/giraph/master/MasterThread.java
@@ -34,6 +34,7 @@ import java.util.Map.Entry; import java.util.TreeMap; +import static org.apache.giraph.conf.GiraphConstants.SPLIT_MASTER_WORKER; import static org.apache.giraph.conf.GiraphConstants.USE_SUPERSTEP_COUNTERS; /** @@ -58,6 +59,8 @@ private final Context context; /** Use superstep counters? */ private final boolean superstepCounterOn; + /** Are master and worker split or not? */ + private final boolean splitMasterWorker; /** Setup seconds */ private double setupSecs = 0d; /** Superstep timer (in seconds) map */ @@ -78,6 +81,7 @@ this.context = context; GiraphTimers.init(context); superstepCounterOn = USE_SUPERSTEP_COUNTERS.get(context.getConfiguration()); + splitMasterWorker = SPLIT_MASTER_WORKER.get(context.getConfiguration()); } /** @@ -118,7 +122,10 @@ while (!superstepState.isExecutionComplete()) { long startSuperstepMillis = System.currentTimeMillis(); long cachedSuperstep = bspServiceMaster.getSuperstep(); - GiraphMetrics.get().resetSuperstepMetrics(cachedSuperstep); + // If master and worker are running together, worker will call reset + if (splitMasterWorker) { + GiraphMetrics.get().resetSuperstepMetrics(cachedSuperstep); + } Class<? extends Computation> computationClass = bspServiceMaster.getMasterCompute().getComputation(); superstepState = bspServiceMaster.coordinateSuperstep();