diff --git a/jps/jps-builders/src/org/jetbrains/jps/incremental/IncProjectBuilder.java b/jps/jps-builders/src/org/jetbrains/jps/incremental/IncProjectBuilder.java index 3506d18cedc7..08e0139c8759 100644 --- a/jps/jps-builders/src/org/jetbrains/jps/incremental/IncProjectBuilder.java +++ b/jps/jps-builders/src/org/jetbrains/jps/incremental/IncProjectBuilder.java @@ -10,7 +10,7 @@ import com.intellij.openapi.util.io.FileUtil; import com.intellij.openapi.util.text.StringUtil; import com.intellij.util.Function; import com.intellij.util.SmartList; -import com.intellij.util.concurrency.AppExecutorUtil; +import com.intellij.util.concurrency.BoundedTaskExecutor; import com.intellij.util.containers.ContainerUtil; import com.intellij.util.containers.MultiMap; import com.intellij.util.io.MappingFailedException; @@ -897,11 +897,17 @@ public class IncProjectBuilder { private final BuildTargetChunk myChunk; private final Set myNotBuiltDependencies = new THashSet<>(); private final List myTasksDependsOnThis = new ArrayList<>(); + private int mySelfScore = 0; + private int myDepsScore = 0; private BuildChunkTask(BuildTargetChunk chunk) { myChunk = chunk; } + private int getScore() { + return myDepsScore + mySelfScore; + } + public BuildTargetChunk getChunk() { return myChunk; } @@ -931,8 +937,12 @@ public class IncProjectBuilder { } private class BuildParallelizer { - private final ExecutorService myParallelBuildExecutor = AppExecutorUtil.createBoundedApplicationPoolExecutor( - "IncProjectBuilder Executor Pool", SharedThreadPool.getInstance(), MAX_BUILDER_THREADS); + private final ExecutorService myParallelBuildExecutor = new BoundedTaskExecutor( + "IncProjectBuilder Executor Pool", SharedThreadPool.getInstance(), MAX_BUILDER_THREADS, true, (o1, o2) -> { + long p1 = o1 instanceof RunnableWithPriority ? ((RunnableWithPriority)o1).priority : 1; + long p2 = o1 instanceof RunnableWithPriority ? ((RunnableWithPriority)o2).priority : 1; + return Long.compare(p2, p1); + }); private final CompileContext myContext; private final BuildProgress myBuildProgress; private final AtomicReference myException = new AtomicReference<>(); @@ -948,15 +958,18 @@ public class IncProjectBuilder { List chunks = targetIndex.getSortedTargetChunks(myContext); myTasks = new ArrayList<>(chunks.size()); - Map, BuildChunkTask> targetToTask = new THashMap<>(); + Map, BuildChunkTask> targetToTask = new THashMap<>(chunks.size()); for (BuildTargetChunk chunk : chunks) { BuildChunkTask task = new BuildChunkTask(chunk); myTasks.add(task); for (BuildTarget target : chunk.getTargets()) { targetToTask.put(target, task); + task.mySelfScore += 1; } } + Map, Collection>> transitiveDependencyCache = new HashMap<>(myTasks.size()); + for (BuildChunkTask task : myTasks) { for (BuildTarget target : task.getChunk().getTargets()) { for (BuildTarget dependency : targetIndex.getDependencies(target, myContext)) { @@ -965,6 +978,12 @@ public class IncProjectBuilder { task.addDependency(depTask); } } + for (BuildTarget dependency : getTransitiveDeps(targetIndex, target, myContext, transitiveDependencyCache)) { + BuildChunkTask depTask = targetToTask.get(dependency); + if (depTask != null && depTask != task) { + depTask.myDepsScore += task.mySelfScore; + } + } } } @@ -997,9 +1016,13 @@ public class IncProjectBuilder { } private void queueTasks(List tasks) { - if (LOG.isDebugEnabled() && !tasks.isEmpty()) { + if (tasks.isEmpty()) return; + ArrayList sorted = new ArrayList<>(tasks); + sorted.sort(Comparator.comparingLong(BuildChunkTask::getScore).reversed()); + + if (LOG.isDebugEnabled()) { final List chunksToLog = new ArrayList<>(); - for (BuildChunkTask task : tasks) { + for (BuildChunkTask task : sorted) { chunksToLog.add(task.getChunk()); } final StringBuilder logBuilder = new StringBuilder("Queuing " + chunksToLog.size() + " chunks in parallel: "); @@ -1009,44 +1032,84 @@ public class IncProjectBuilder { } LOG.debug(logBuilder.toString()); } - for (BuildChunkTask task : tasks) { + for (BuildChunkTask task : sorted) { queueTask(task); } } + private abstract class RunnableWithPriority implements Runnable { + public final int priority; + + RunnableWithPriority(int priority) { + this.priority = priority; + } + } + private void queueTask(final BuildChunkTask task) { final CompileContext chunkLocalContext = createContextWrapper(myContext); - myParallelBuildExecutor.execute(() -> { - try { + myParallelBuildExecutor.execute(new RunnableWithPriority(task.getScore()) { + @Override + public void run() { try { - if (myException.get() == null) { - buildChunkIfAffected(chunkLocalContext, myContext.getScope(), task.getChunk(), myBuildProgress); + try { + if (myException.get() == null) { + buildChunkIfAffected(chunkLocalContext, myContext.getScope(), task.getChunk(), myBuildProgress); + } + } + finally { + myProjectDescriptor.dataManager.closeSourceToOutputStorages(Collections.singletonList(task.getChunk())); + myProjectDescriptor.dataManager.flush(true); } } + catch (Throwable e) { + myException.compareAndSet(null, e); + LOG.info(e); + } finally { - myProjectDescriptor.dataManager.closeSourceToOutputStorages(Collections.singletonList(task.getChunk())); - myProjectDescriptor.dataManager.flush(true); - } - } - catch (Throwable e) { - myException.compareAndSet(null, e); - LOG.info(e); - } - finally { - LOG.debug("Finished compilation of " + task.getChunk().toString()); - myTasksCountDown.countDown(); - List nextTasks; - synchronized (myQueueLock) { - nextTasks = task.markAsFinishedAndGetNextReadyTasks(); - } - if (!nextTasks.isEmpty()) { - queueTasks(nextTasks); + LOG.debug("Finished compilation of " + task.getChunk().toString()); + myTasksCountDown.countDown(); + List nextTasks; + synchronized (myQueueLock) { + nextTasks = task.markAsFinishedAndGetNextReadyTasks(); + } + if (!nextTasks.isEmpty()) { + queueTasks(nextTasks); + } } } }); } } + private static Iterable> getTransitiveDeps(BuildTargetIndex index, + BuildTarget target, + CompileContext context, + Map, Collection>> cache) { + if (cache.containsKey(target)) { + return cache.get(target); + } + Set> result = new HashSet<>(); + LinkedList> queue = new LinkedList<>(); + queue.add(target); + result.add(target); + while (!queue.isEmpty()) { + BuildTarget next = queue.pop(); + Collection> transitive = cache.get(next); + if (transitive != null) { + result.addAll(transitive); + } + else { + Collection> dependencies = index.getDependencies(next, context); + for (BuildTarget dependency : dependencies) { + if (dependency != target && result.add(dependency)) queue.add(dependency); + } + } + } + result.remove(target); + cache.put(target, result); + return result; + } + private void buildChunkIfAffected(CompileContext context, CompileScope scope, BuildTargetChunk chunk, BuildProgress buildProgress) throws ProjectBuildException { if (isAffected(scope, chunk)) { diff --git a/platform/util/src/com/intellij/util/concurrency/BoundedTaskExecutor.java b/platform/util/src/com/intellij/util/concurrency/BoundedTaskExecutor.java index 6099d74b1397..5523a50596b3 100644 --- a/platform/util/src/com/intellij/util/concurrency/BoundedTaskExecutor.java +++ b/platform/util/src/com/intellij/util/concurrency/BoundedTaskExecutor.java @@ -1,4 +1,4 @@ -// Copyright 2000-2019 JetBrains s.r.o. Use of this source code is governed by the Apache 2.0 license that can be found in the LICENSE file. +// Copyright 2000-2020 JetBrains s.r.o. Use of this source code is governed by the Apache 2.0 license that can be found in the LICENSE file. package com.intellij.util.concurrency; import com.intellij.diagnostic.StartUpMeasurer; @@ -14,9 +14,11 @@ import com.intellij.util.ReflectionUtil; import com.intellij.util.containers.ContainerUtil; import org.jetbrains.annotations.Async; import org.jetbrains.annotations.NotNull; +import org.jetbrains.annotations.Nullable; import java.util.ArrayList; import java.util.Collections; +import java.util.Comparator; import java.util.List; import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicLong; @@ -38,11 +40,15 @@ public final class BoundedTaskExecutor extends AbstractExecutorService { // low 32 bits: number of tasks running (or trying to run) // high 32 bits: myTaskQueue modification stamp private final AtomicLong myStatus = new AtomicLong(); - private final BlockingQueue myTaskQueue = new LinkedBlockingQueue<>(); + private final BlockingQueue myTaskQueue; private final boolean myChangeThreadName; BoundedTaskExecutor(@NotNull String name, @NotNull Executor backendExecutor, int maxThreads, boolean changeThreadName) { + this(name, backendExecutor, maxThreads, changeThreadName, null); + } + + public BoundedTaskExecutor(@NotNull String name, @NotNull Executor backendExecutor, int maxThreads, boolean changeThreadName, @Nullable Comparator comparator) { myName = StringUtil.capitalize(name); myBackendExecutor = backendExecutor; if (maxThreads < 1) { @@ -53,6 +59,11 @@ public final class BoundedTaskExecutor extends AbstractExecutorService { } myMaxThreads = maxThreads; myChangeThreadName = changeThreadName; + if (comparator != null) { + myTaskQueue = new PriorityBlockingQueue<>(11, comparator); + } else { + myTaskQueue = new LinkedBlockingQueue<>(); + } } /** @deprecated use {@link AppExecutorUtil#createBoundedApplicationPoolExecutor(String, Executor, int)} instead */