1. generalized SequentialTaskExecutor: now supports specified number of simultaneous threads

2. if parallel build enabled, force in-memory temp caches storage
This commit is contained in:
Eugene Zhuravlev
2012-09-02 16:57:09 +02:00
parent abe476f2c4
commit ac2b9c0ffa
6 changed files with 106 additions and 105 deletions
@@ -0,0 +1,83 @@
package org.jetbrains.jps.api;
import org.jetbrains.annotations.Nullable;
import java.util.Queue;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
/**
* @author Eugene Zhuravlev
* Date: 9/2/12
*/
public class BoundedTaskExecutor implements Executor {
protected final Executor myBackendExecutor;
private final int myMaxTasks;
private final AtomicInteger myInProgress = new AtomicInteger(0);
private final Queue<FutureTask> myTaskQueue = new LinkedBlockingQueue<FutureTask>();
private final Runnable USER_TASK_RUNNER = new Runnable() {
public void run() {
final FutureTask task = myTaskQueue.poll();
try {
if (task != null && !task.isCancelled()) {
task.run();
}
}
finally {
myInProgress.decrementAndGet();
if (!myTaskQueue.isEmpty()) {
processQueue();
}
}
}
};
public BoundedTaskExecutor(Executor backendExecutor, int maxSimultaneousTasks) {
myBackendExecutor = backendExecutor;
myMaxTasks = Math.max(maxSimultaneousTasks, 1);
}
@Override
public void execute(Runnable task) {
submit(task);
}
public Future submit(Runnable task) {
final RunnableFuture<Void> future = queueTask(new FutureTask<Void>(task, null));
if (future == null) {
throw new RuntimeException("Failed to queue task: " + task);
}
return future;
}
public <T> Future<T> submit(Callable<T> task) {
final RunnableFuture<T> future = queueTask(new FutureTask<T>(task));
if (future == null) {
throw new RuntimeException("Failed to queue task: " + task);
}
return future;
}
@Nullable
private <T> RunnableFuture<T> queueTask(FutureTask<T> futureTask) {
if (myTaskQueue.offer(futureTask)) {
processQueue();
return futureTask;
}
return null;
}
protected void processQueue() {
while (true) {
final int count = myInProgress.get();
if (count >= myMaxTasks) {
return;
}
if (myInProgress.compareAndSet(count, count + 1)) {
break;
}
}
myBackendExecutor.execute(USER_TASK_RUNNER);
}
}
@@ -1,61 +1,14 @@
package org.jetbrains.jps.api;
import java.util.Queue;
import java.util.concurrent.Executor;
import java.util.concurrent.FutureTask;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.RunnableFuture;
import java.util.concurrent.atomic.AtomicBoolean;
/**
* @author Eugene Zhuravlev
* Date: 9/24/11
*/
public class SequentialTaskExecutor implements Executor {
private final Executor myExecutor;
private final Queue<FutureTask> myTaskQueue = new LinkedBlockingQueue<FutureTask>();
private final AtomicBoolean myInProgress = new AtomicBoolean(false);
private final Runnable USER_TASK_RUNNER = new Runnable() {
public void run() {
final FutureTask task = myTaskQueue.poll();
try {
if (task != null && !task.isCancelled()) {
task.run();
}
}
finally {
myInProgress.set(false);
if (!myTaskQueue.isEmpty()) {
processQueue();
}
}
}
};
public class SequentialTaskExecutor extends BoundedTaskExecutor {
public SequentialTaskExecutor(Executor executor) {
myExecutor = executor;
super(executor, 1);
}
@Override
public void execute(Runnable task) {
submit(task);
}
public RunnableFuture submit(Runnable task) {
final FutureTask futureTask = new FutureTask(task, null);
if (myTaskQueue.offer(futureTask)) {
processQueue();
}
else {
throw new RuntimeException("Failed to queue task: " + task);
}
return futureTask;
}
protected void processQueue() {
if (!myInProgress.getAndSet(true)) {
myExecutor.execute(USER_TASK_RUNNER);
}
}
}
@@ -1,47 +0,0 @@
package org.jetbrains.jps.api;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
/**
* @author Eugene Zhuravlev
* Date: 3/29/12
*/
public class SharedBuilderThreadPool {
private static final int MAX_BUILDER_THREADS;
static {
int maxThreads = 4;
try {
maxThreads = Math.max(2, Integer.parseInt(System.getProperty(GlobalOptions.COMPILE_PARALLEL_MAX_THREADS_OPTION, "4")));
}
catch (NumberFormatException ignored) {
}
MAX_BUILDER_THREADS = maxThreads;
}
private static final ExecutorService ourBuilderPool = Executors.newFixedThreadPool(Math.min(MAX_BUILDER_THREADS, Math.max(2, Runtime.getRuntime().availableProcessors())));
public static final SharedBuilderThreadPool INSTANCE = new SharedBuilderThreadPool();
private SharedBuilderThreadPool() {
}
/** @noinspection MethodMayBeStatic*/
public Future<?> submitBuildTask(final Runnable task) {
return _submit(task, ourBuilderPool);
}
private static Future<?> _submit(final Runnable task, final ExecutorService service) {
return service.submit(new Runnable() {
public void run() {
try {
task.run();
}
finally {
Thread.interrupted(); // reset interrupted status
}
}
});
}
}
@@ -32,6 +32,8 @@ import java.util.*;
*/
public class BuildRunner {
private static final Logger LOG = Logger.getInstance("#org.jetbrains.jps.cmdline.BuildRunner");
public static final boolean PARALLEL_BUILD_ENABLED = Boolean.parseBoolean(System.getProperty(GlobalOptions.COMPILE_PARALLEL_OPTION, "false"));
private static final boolean STORE_TEMP_CACHES_IN_MEMORY = PARALLEL_BUILD_ENABLED || System.getProperty(GlobalOptions.USE_MEMORY_TEMP_CACHE_OPTION) != null;
private final JpsModelLoader myModelLoader;
private final Set<String> myModules;
private final List<String> myArtifacts;
@@ -50,12 +52,11 @@ public class BuildRunner {
}
public ProjectDescriptor load(MessageHandler msgHandler, File dataStorageRoot, BuildFSState fsState) throws IOException {
final boolean inMemoryMappingsDelta = System.getProperty(GlobalOptions.USE_MEMORY_TEMP_CACHE_OPTION) != null;
ProjectTimestamps projectTimestamps = null;
BuildDataManager dataManager = null;
try {
projectTimestamps = new ProjectTimestamps(dataStorageRoot);
dataManager = new BuildDataManager(dataStorageRoot, inMemoryMappingsDelta);
dataManager = new BuildDataManager(dataStorageRoot, STORE_TEMP_CACHES_IN_MEMORY);
if (dataManager.versionDiffers()) {
myForceCleanCaches = true;
msgHandler.processMessage(new CompilerMessage("build", BuildMessage.Kind.INFO, "Dependency data format has changed, project rebuild required"));
@@ -73,7 +74,7 @@ public class BuildRunner {
myForceCleanCaches = true;
FileUtil.delete(dataStorageRoot);
projectTimestamps = new ProjectTimestamps(dataStorageRoot);
dataManager = new BuildDataManager(dataStorageRoot, inMemoryMappingsDelta);
dataManager = new BuildDataManager(dataStorageRoot, STORE_TEMP_CACHES_IN_MEMORY);
// second attempt succeded
msgHandler.processMessage(new CompilerMessage("build", BuildMessage.Kind.INFO, "Project rebuild forced: " + e.getMessage()));
}
@@ -11,10 +11,11 @@ import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.ether.dependencyView.Callbacks;
import org.jetbrains.jps.*;
import org.jetbrains.jps.api.BoundedTaskExecutor;
import org.jetbrains.jps.api.CanceledStatus;
import org.jetbrains.jps.api.GlobalOptions;
import org.jetbrains.jps.api.RequestFuture;
import org.jetbrains.jps.api.SharedBuilderThreadPool;
import org.jetbrains.jps.cmdline.BuildRunner;
import org.jetbrains.jps.cmdline.ProjectDescriptor;
import org.jetbrains.jps.incremental.fs.BuildFSState;
import org.jetbrains.jps.incremental.fs.RootDescriptor;
@@ -50,7 +51,17 @@ public class IncProjectBuilder {
public static final String BUILD_NAME = "EXTERNAL BUILD";
private static final String CLASSPATH_INDEX_FINE_NAME = "classpath.index";
private static final boolean GENERATE_CLASSPATH_INDEX = Boolean.parseBoolean(System.getProperty(GlobalOptions.GENERATE_CLASSPATH_INDEX_OPTION, "false"));
private static final boolean PARALLEL_BUILD_ENABLED = Boolean.parseBoolean(System.getProperty(GlobalOptions.COMPILE_PARALLEL_OPTION, "false"));
private static final int MAX_BUILDER_THREADS;
static {
int maxThreads = 4;
try {
maxThreads = Math.max(2, Integer.parseInt(System.getProperty(GlobalOptions.COMPILE_PARALLEL_MAX_THREADS_OPTION, "4")));
}
catch (NumberFormatException ignored) {
}
MAX_BUILDER_THREADS = maxThreads;
}
private final BoundedTaskExecutor myParallelBuildExecutor = new BoundedTaskExecutor(SharedThreadPool.getInstance(), Math.min(MAX_BUILDER_THREADS, Math.max(2, Runtime.getRuntime().availableProcessors())));
private final ProjectDescriptor myProjectDescriptor;
private final BuilderRegistry myBuilderRegistry;
@@ -191,7 +202,7 @@ public class IncProjectBuilder {
private void runBuild(CompileContextImpl context, boolean forceCleanCaches) throws ProjectBuildException {
context.setDone(0.0f);
LOG.info("Building project '" + context.getProjectDescriptor().project.getProjectName() + "'; isRebuild:" + context.isProjectRebuild() + "; isMake:" + context.isMake() + " parallel compilation:" + PARALLEL_BUILD_ENABLED);
LOG.info("Building project '" + context.getProjectDescriptor().project.getProjectName() + "'; isRebuild:" + context.isProjectRebuild() + "; isMake:" + context.isMake() + " parallel compilation:" + BuildRunner.PARALLEL_BUILD_ENABLED);
for (ProjectLevelBuilder builder : myBuilderRegistry.getProjectLevelBuilders()) {
builder.buildStarted(context);
@@ -405,7 +416,7 @@ public class IncProjectBuilder {
final CompileScope scope = context.getScope();
final ProjectDescriptor pd = context.getProjectDescriptor();
try {
if (PARALLEL_BUILD_ENABLED) {
if (BuildRunner.PARALLEL_BUILD_ENABLED) {
final List<ChunkGroup> chunkGroups = buildChunkGroups(context, chunks);
for (ChunkGroup group : chunkGroups) {
final List<ModuleChunk> groupChunks = group.getChunks();
@@ -431,7 +442,7 @@ public class IncProjectBuilder {
for (final ModuleChunk chunk : groupChunks) {
final CompileContext chunkLocalContext = createContextWrapper(context);
SharedBuilderThreadPool.INSTANCE.submitBuildTask(new Runnable() {
myParallelBuildExecutor.execute(new Runnable() {
@Override
public void run() {
try {
@@ -164,7 +164,7 @@ compiler.process.vm.options.description=Additional options for the build process
compiler.process.use.memory.temp.cache=true
# suppress inspection "UnusedProperty"
compiler.process.use.memory.temp.cache.description=Store temporary data in memory for faster compilation; requires larger heap size for the build process.
compiler.process.use.memory.temp.cache.description=Store temporary data in memory for faster compilation; requires larger heap size for the build process. If parallel build is enabled, the option is ignored and temp data is always stored in memory.
compiler.process.use.external.javac=false
# suppress inspection "UnusedProperty"