IDEA-311941 Employ SimpleMergingQueue for ProjectStructureDaemonAnalyzer

ProjectStructureDaemonAnalyzer often gets all modules of the current project
for update at once. If there are quite a lot of them (e.g., 5k), MergingUpdateQueue
wastes a huge amount of time trying to merge them one by one with complicated rules,
though the rules are extremely simple in our case: it's enough to just check updates
for equality. A simple hash map works great for this task.
So, I added SimpleMergingQueue.

This fix helps to reduce queueing tasks for intellij monorepo from ~6sec to less than
0.1sec.

GitOrigin-RevId: 33eaf1e7ccc66a1edb28addf7bbaedd00444a30b
This commit is contained in:
Max Medvedev
2024-01-31 13:05:42 +00:00
committed by intellij-monorepo-bot
parent 92269baf0e
commit a4cf21904a
4 changed files with 292 additions and 59 deletions
@@ -9,7 +9,7 @@ import com.intellij.openapi.util.Disposer;
import com.intellij.openapi.util.MultiValuesMap;
import com.intellij.util.Alarm;
import com.intellij.util.EventDispatcher;
import com.intellij.util.ui.update.MergingUpdateQueue;
import com.intellij.util.containers.ContainerUtil;
import com.intellij.util.ui.update.Update;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
@@ -25,8 +25,8 @@ public class ProjectStructureDaemonAnalyzer implements Disposable {
private final Set<ProjectStructureElement> myElementWithNotCalculatedUsages = new HashSet<>();
private final Set<ProjectStructureElement> myElementsToShowWarningIfUnused = new HashSet<>();
private final Map<ProjectStructureElement, ProjectStructureProblemDescription> myWarningsAboutUnused = new HashMap<>();
private final MergingUpdateQueue myAnalyzerQueue;
private final MergingUpdateQueue myResultsUpdateQueue;
private final SimpleMergingQueue<Runnable> myAnalyzerQueue;
private final SimpleMergingQueue<Runnable> myResultsUpdateQueue;
private final EventDispatcher<ProjectStructureDaemonAnalyzerListener> myDispatcher = EventDispatcher.create(ProjectStructureDaemonAnalyzerListener.class);
private final AtomicBoolean myStopped = new AtomicBoolean(false);
private final ProjectConfigurationProblems myProjectConfigurationProblems;
@@ -34,9 +34,9 @@ public class ProjectStructureDaemonAnalyzer implements Disposable {
public ProjectStructureDaemonAnalyzer(@NotNull StructureConfigurableContext context) {
Disposer.register(context, this);
myProjectConfigurationProblems = new ProjectConfigurationProblems(this, context);
myAnalyzerQueue = new MergingUpdateQueue("Project Structure Daemon Analyzer", 300, false, null, this, null, Alarm.ThreadToUse.POOLED_THREAD);
myResultsUpdateQueue = new MergingUpdateQueue("Project Structure Analysis Results Updater", 300, false, MergingUpdateQueue.ANY_COMPONENT,
this, null, Alarm.ThreadToUse.SWING_THREAD);
myAnalyzerQueue = new SimpleMergingQueue<>("Project Structure Daemon Analyzer", 300, false, Alarm.ThreadToUse.POOLED_THREAD, this);
myResultsUpdateQueue = new SimpleMergingQueue<>("Project Structure Analysis Results Updater", 300, false, Alarm.ThreadToUse.SWING_THREAD, this);
}
private void doUpdate(@NotNull ProjectStructureElement element) {
@@ -168,11 +168,9 @@ public class ProjectStructureDaemonAnalyzer implements Disposable {
public void stop() {
LOG.debug("analyzer stopped");
myStopped.set(true);
myAnalyzerQueue.cancelAllUpdates();
myResultsUpdateQueue.cancelAllUpdates();
clearCaches();
myAnalyzerQueue.deactivate();
myResultsUpdateQueue.deactivate();
myAnalyzerQueue.stop();
myResultsUpdateQueue.stop();
}
public void clearCaches() {
@@ -189,6 +187,7 @@ public class ProjectStructureDaemonAnalyzer implements Disposable {
}
myProblemHolders.clear();
LOG.debug("Adding to queue updates for " + toUpdate.size() + " problematic elements");
for (ProjectStructureElement element : toUpdate) {
queueUpdate(element);
}
@@ -197,8 +196,8 @@ public class ProjectStructureDaemonAnalyzer implements Disposable {
@Override
public void dispose() {
myStopped.set(true);
myAnalyzerQueue.cancelAllUpdates();
myResultsUpdateQueue.cancelAllUpdates();
myAnalyzerQueue.stop();
myResultsUpdateQueue.stop();
}
@Nullable
@@ -222,8 +221,8 @@ public class ProjectStructureDaemonAnalyzer implements Disposable {
public void reset() {
LOG.debug("analyzer started");
myAnalyzerQueue.activate();
myResultsUpdateQueue.activate();
myAnalyzerQueue.start();
myResultsUpdateQueue.start();
myAnalyzerQueue.queue(new Update("reset") {
@Override
public void run() {
@@ -241,25 +240,11 @@ public class ProjectStructureDaemonAnalyzer implements Disposable {
myProjectConfigurationProblems.clearProblems();
}
private class AnalyzeElementUpdate extends Update {
private final class AnalyzeElementUpdate implements Runnable {
private final @NotNull ProjectStructureElement myElement;
private final @NotNull Object @NotNull [] myEqualityObjects;
AnalyzeElementUpdate(@NotNull ProjectStructureElement element) {
super(element);
myElement = element;
myEqualityObjects = new Object[]{myElement};
}
@Override
public boolean canEat(@NotNull Update update) {
if (!(update instanceof AnalyzeElementUpdate other)) return false;
return myElement.equals(other.myElement);
}
@Override
public Object @NotNull [] getEqualityObjects() {
return myEqualityObjects;
}
@Override
@@ -271,23 +256,28 @@ public class ProjectStructureDaemonAnalyzer implements Disposable {
LOG.error(t);
}
}
}
private class UsagesCollectedUpdate extends Update {
private final @NotNull ProjectStructureElement myElement;
private final @NotNull List<? extends ProjectStructureElementUsage> myUsages;
private final @NotNull Object @NotNull [] myEqualityObjects;
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (!(o instanceof AnalyzeElementUpdate update)) return false;
UsagesCollectedUpdate(@NotNull ProjectStructureElement element, @NotNull List<? extends ProjectStructureElementUsage> usages) {
super(element);
myElement = element;
myUsages = usages;
myEqualityObjects = new Object[]{element, "usages collected"};
return myElement.equals(update.myElement);
}
@Override
public Object @NotNull [] getEqualityObjects() {
return myEqualityObjects;
public int hashCode() {
return myElement.hashCode();
}
}
private final class UsagesCollectedUpdate implements Runnable {
private final @NotNull ProjectStructureElement myElement;
private final @NotNull List<? extends ProjectStructureElementUsage> myUsages;
UsagesCollectedUpdate(@NotNull ProjectStructureElement element, @NotNull List<? extends ProjectStructureElementUsage> usages) {
myElement = element;
myUsages = usages;
}
@Override
@@ -299,23 +289,28 @@ public class ProjectStructureDaemonAnalyzer implements Disposable {
}
updateUsages(myElement, myUsages);
}
}
private class ProblemsComputedUpdate extends Update {
private final @NotNull ProjectStructureElement myElement;
private final @NotNull ProjectStructureProblemsHolderImpl myProblemsHolder;
private final @NotNull Object @NotNull [] myEqualityObjects;
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (!(o instanceof UsagesCollectedUpdate other)) return false;
ProblemsComputedUpdate(@NotNull ProjectStructureElement element, @NotNull ProjectStructureProblemsHolderImpl problemsHolder) {
super(element);
myElement = element;
myProblemsHolder = problemsHolder;
myEqualityObjects = new Object[]{element, "problems computed"};
return myElement.equals(other.myElement);
}
@Override
public Object @NotNull [] getEqualityObjects() {
return myEqualityObjects;
public int hashCode() {
return myElement.hashCode() + 1;
}
}
private final class ProblemsComputedUpdate implements Runnable {
private final @NotNull ProjectStructureElement myElement;
private final @NotNull ProjectStructureProblemsHolderImpl myProblemsHolder;
ProblemsComputedUpdate(@NotNull ProjectStructureElement element, @NotNull ProjectStructureProblemsHolderImpl problemsHolder) {
myElement = element;
myProblemsHolder = problemsHolder;
}
@Override
@@ -332,16 +327,35 @@ public class ProjectStructureDaemonAnalyzer implements Disposable {
myProblemHolders.put(myElement, myProblemsHolder);
myDispatcher.getMulticaster().problemsChanged(myElement);
}
}
private final class ReportUnusedElementsUpdate extends Update {
private ReportUnusedElementsUpdate() {
super("unused elements");
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (!(o instanceof ProblemsComputedUpdate update)) return false;
return myElement.equals(update.myElement);
}
@Override
public int hashCode() {
return myElement.hashCode() + 2;
}
}
private final class ReportUnusedElementsUpdate implements Runnable {
@Override
public void run() {
reportUnusedElements();
}
@Override
public int hashCode() {
return 3;
}
@Override
public boolean equals(Object obj) {
return obj instanceof ReportUnusedElementsUpdate;
}
}
}
@@ -0,0 +1,101 @@
// Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
package com.intellij.openapi.roots.ui.configuration.projectRoot.daemon
import com.intellij.openapi.Disposable
import com.intellij.openapi.diagnostic.logger
import com.intellij.openapi.util.Disposer
import com.intellij.util.Alarm
import com.intellij.util.Alarm.ThreadToUse
import com.intellij.util.ObjectUtils
import org.jetbrains.annotations.ApiStatus.Internal
import org.jetbrains.annotations.VisibleForTesting
/**
* When you queue a task, the queue waits for [delayMillis] and then runs the task.
* If a new equal task comes during [delayMillis], the new task is ignored.
* The tasks are executed in the same order they were queued.
*
* The tasks must implement proper equals&hashcode.
*/
@Internal
@VisibleForTesting
class SimpleMergingQueue<T : Runnable>(
private val name: String,
private val delayMillis: Int,
@Volatile private var isActive: Boolean,
threadToUse: ThreadToUse,
disposableParent: Disposable
) : Disposable {
private var taskHolder = LinkedHashSet<T>()
private val timer = Alarm(threadToUse, this)
private val lock = ObjectUtils.sentinel("SimpleMergingQueue($name)")
init {
Disposer.register(disposableParent, this)
}
fun start() {
synchronized(lock) {
if (isActive) return
isActive = true
if (taskHolder.isNotEmpty()) {
waitAndRun()
}
}
}
fun stop() {
synchronized(lock) {
isActive = false
timer.cancelAllRequests()
}
}
fun queue(task: T) {
queue(listOf(task))
}
fun queue(tasks: List<T>) {
synchronized(lock) {
taskHolder.addAll(tasks)
if (isActive) {
waitAndRun()
}
}
}
private fun waitAndRun() {
timer.cancelAllRequests()
timer.addRequest({ flush() }, delayMillis)
}
private fun flush() {
val tasks: Set<T>
synchronized(taskHolder) {
if (!isActive) return
tasks = taskHolder
taskHolder = LinkedHashSet()
}
for (task in tasks) {
try {
task.run()
}
catch (e: Throwable) {
log.error(e)
}
}
}
override fun toString(): String =
"SimpleMergingQueue(name='$name', delayMillis=$delayMillis, isActive=$isActive)"
override fun dispose() {
synchronized(lock) {
stop()
taskHolder = LinkedHashSet()
}
}
}
private val log = logger<SimpleMergingQueue<*>>()
@@ -0,0 +1,118 @@
// Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
package com.intellij.openapi.roots.ui.configuration.projectRoot.daemon
import com.intellij.openapi.Disposable
import com.intellij.openapi.util.Disposer
import com.intellij.testFramework.LoggedErrorProcessor
import com.intellij.testFramework.UsefulTestCase.assertOrderedEquals
import com.intellij.util.Alarm.ThreadToUse
import org.junit.jupiter.api.AfterEach
import org.junit.jupiter.api.Assertions
import org.junit.jupiter.api.BeforeEach
import org.junit.jupiter.params.ParameterizedTest
import org.junit.jupiter.params.provider.EnumSource
import java.util.concurrent.Semaphore
class SimpleMergeQueueTest {
private lateinit var disposable: Disposable
@BeforeEach
fun setUp() {
disposable = Disposer.newDisposable()
}
@AfterEach
fun tearDown() {
Disposer.dispose(disposable)
}
@ParameterizedTest
@EnumSource(ThreadToUse::class, names = ["SWING_THREAD", "POOLED_THREAD"])
fun duplicateIsRemoved(threadToUse: ThreadToUse) {
val list = mutableListOf<String>()
executeInQueue(threadToUse) { queue ->
queue.queue(Task("0", "0", list))
queue.queue(Task("1", "1", list))
queue.queue(Task("2", "2", list))
queue.queue(Task("0", "3", list)) // equals to the first task
}
synchronized(list) {
assertOrderedEquals(list, "0", "1", "2")
}
}
@ParameterizedTest
@EnumSource(ThreadToUse::class, names = ["SWING_THREAD", "POOLED_THREAD"])
fun duplicateIsNotRemovedAfterTimeout() {
val list = mutableListOf<String>()
val queue = SimpleMergingQueue<Runnable>("test", 100, false, ThreadToUse.POOLED_THREAD, disposable)
queue.queue(Task("0", "0", list))
queue.queue(Task("1", "1", list))
waitForQueue(queue)
queue.queue(Task("0", "0", list))
queue.queue(Task("1", "1", list))
waitForQueue(queue)
synchronized(list) {
assertOrderedEquals(list, "0", "1", "0", "1")
}
}
@ParameterizedTest
@EnumSource(ThreadToUse::class, names = ["SWING_THREAD", "POOLED_THREAD"])
fun errorInTaskDoesNotPreventConsequentTasksFromExecuting() {
val list = mutableListOf<String>()
val error = RuntimeException()
// testLogger usually throws. Instead, let's intercept the error. In this case, the test works the same way as in production
val loggerError = LoggedErrorProcessor.executeAndReturnLoggedError {
executeInQueue(ThreadToUse.POOLED_THREAD) { queue ->
queue.queue(Task("0", "0", list))
queue.queue {
throw error // this error should not prevent the next task from execution
}
queue.queue(Task("1", "1", list))
}
}
Assertions.assertEquals(loggerError, error)
synchronized(list) {
assertOrderedEquals(list, "0", "1")
}
}
private fun waitForQueue(queue: SimpleMergingQueue<Runnable>) {
val semaphore = Semaphore(0)
queue.queue(Runnable {
semaphore.release()
})
queue.start()
semaphore.acquire() // wait for tasks to be processed
}
private fun executeInQueue(threadToUse: ThreadToUse, block: (SimpleMergingQueue<Runnable>) -> Unit) {
val queue = SimpleMergingQueue<Runnable>("test", 100, false, threadToUse, disposable)
block(queue)
waitForQueue(queue)
}
private class Task(val equality: String, val text: String, val list: MutableList<String>) : Runnable {
override fun run() {
synchronized(list) {
list += text
}
}
override fun hashCode(): Int {
return equality.hashCode()
}
override fun equals(other: Any?): Boolean {
return other is Task && equality == other.equality
}
}
}
@@ -1,4 +1,4 @@
// Copyright 2000-2023 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
// Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
package com.intellij.testFramework;
import com.intellij.util.ThrowableRunnable;
@@ -11,7 +11,7 @@ import java.util.Set;
import java.util.concurrent.atomic.AtomicReference;
public class LoggedErrorProcessor {
private static LoggedErrorProcessor ourInstance = new LoggedErrorProcessor();
private static volatile LoggedErrorProcessor ourInstance = new LoggedErrorProcessor();
static @NotNull LoggedErrorProcessor getInstance() {
return ourInstance;