This is an automated email from the ASF dual-hosted git repository.
btellier pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/james-project.git
The following commit(s) were added to refs/heads/master by this push:
new 3bef1d2 JAMES-3652 Avoid locking in protocol task execution (#666)
3bef1d2 is described below
commit 3bef1d22e7891a82076e54a140c111bf5ba53ecf
Author: Benoit TELLIER <[email protected]>
AuthorDate: Tue Sep 28 09:29:05 2021 +0700
JAMES-3652 Avoid locking in protocol task execution (#666)
---
.../JMXEnabledScheduledThreadPoolExecutor.java | 21 ++++++++-------------
.../concurrent/JMXEnabledThreadPoolExecutor.java | 22 ++++++++--------------
.../lib/AbstractStateCompositeProcessor.java | 6 +++---
.../lib/AbstractStateMailetProcessor.java | 6 +++---
...nabledOrderedMemoryAwareThreadPoolExecutor.java | 21 ++++++++-------------
5 files changed, 30 insertions(+), 46 deletions(-)
diff --git
a/server/container/util/src/main/java/org/apache/james/util/concurrent/JMXEnabledScheduledThreadPoolExecutor.java
b/server/container/util/src/main/java/org/apache/james/util/concurrent/JMXEnabledScheduledThreadPoolExecutor.java
index 0357aab..331a9ff 100644
---
a/server/container/util/src/main/java/org/apache/james/util/concurrent/JMXEnabledScheduledThreadPoolExecutor.java
+++
b/server/container/util/src/main/java/org/apache/james/util/concurrent/JMXEnabledScheduledThreadPoolExecutor.java
@@ -19,10 +19,10 @@
package org.apache.james.util.concurrent;
import java.lang.management.ManagementFactory;
-import java.util.ArrayList;
-import java.util.Collections;
import java.util.List;
import java.util.concurrent.ScheduledThreadPoolExecutor;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
import javax.management.MBeanServer;
import javax.management.ObjectName;
@@ -35,10 +35,9 @@ import javax.management.ObjectName;
public class JMXEnabledScheduledThreadPoolExecutor extends
ScheduledThreadPoolExecutor implements
JMXEnabledScheduledThreadPoolExecutorMBean {
private final String jmxPath;
- private final List<Runnable> inProgress = Collections.synchronizedList(new
ArrayList<>());
private final ThreadLocal<Long> startTime = new ThreadLocal<>();
- private long totalTime;
- private int totalTasks;
+ private final AtomicLong totalTime = new AtomicLong(0);
+ private final AtomicInteger totalTasks = new AtomicInteger(0);
private MBeanServer mbeanServer;
private String mbeanName;
@@ -59,18 +58,14 @@ public class JMXEnabledScheduledThreadPoolExecutor extends
ScheduledThreadPoolEx
@Override
protected void beforeExecute(Thread t, Runnable r) {
super.beforeExecute(t, r);
- inProgress.add(r);
startTime.set(System.currentTimeMillis());
}
@Override
protected void afterExecute(Runnable r, Throwable t) {
long time = System.currentTimeMillis() - startTime.get();
- synchronized (this) {
- totalTime += time;
- ++totalTasks;
- }
- inProgress.remove(r);
+ totalTasks.incrementAndGet();
+ totalTime.addAndGet(time);
super.afterExecute(r, t);
}
@@ -121,12 +116,12 @@ public class JMXEnabledScheduledThreadPoolExecutor
extends ScheduledThreadPoolEx
@Override
public synchronized int getTotalTasks() {
- return totalTasks;
+ return totalTasks.get();
}
@Override
public synchronized double getAverageTaskTime() {
- return (totalTasks == 0) ? 0 : totalTime / totalTasks;
+ return (totalTasks.get() == 0) ? 0 : totalTime.get() /
totalTasks.get();
}
@Override
diff --git
a/server/container/util/src/main/java/org/apache/james/util/concurrent/JMXEnabledThreadPoolExecutor.java
b/server/container/util/src/main/java/org/apache/james/util/concurrent/JMXEnabledThreadPoolExecutor.java
index b6e0051..9dcf209 100644
---
a/server/container/util/src/main/java/org/apache/james/util/concurrent/JMXEnabledThreadPoolExecutor.java
+++
b/server/container/util/src/main/java/org/apache/james/util/concurrent/JMXEnabledThreadPoolExecutor.java
@@ -19,14 +19,14 @@
package org.apache.james.util.concurrent;
import java.lang.management.ManagementFactory;
-import java.util.ArrayList;
-import java.util.Collections;
import java.util.List;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.SynchronousQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
import javax.management.MBeanServer;
import javax.management.ObjectName;
@@ -37,10 +37,9 @@ import javax.management.ObjectName;
public class JMXEnabledThreadPoolExecutor extends ThreadPoolExecutor
implements JMXEnabledThreadPoolExecutorMBean {
private final String jmxPath;
- private final List<Runnable> inProgress = Collections.synchronizedList(new
ArrayList<>());
private final ThreadLocal<Long> startTime = new ThreadLocal<>();
- private long totalTime;
- private int totalTasks;
+ private final AtomicLong totalTime = new AtomicLong(0);
+ private final AtomicInteger totalTasks = new AtomicInteger(0);
private MBeanServer mbeanServer;
private String mbeanName;
@@ -53,18 +52,14 @@ public class JMXEnabledThreadPoolExecutor extends
ThreadPoolExecutor implements
@Override
protected void beforeExecute(Thread t, Runnable r) {
super.beforeExecute(t, r);
- inProgress.add(r);
startTime.set(System.currentTimeMillis());
}
@Override
protected void afterExecute(Runnable r, Throwable t) {
long time = System.currentTimeMillis() - startTime.get();
- synchronized (this) {
- totalTime += time;
- ++totalTasks;
- }
- inProgress.remove(r);
+ totalTasks.incrementAndGet();
+ totalTime.addAndGet(time);
super.afterExecute(r, t);
}
@@ -84,7 +79,6 @@ public class JMXEnabledThreadPoolExecutor extends
ThreadPoolExecutor implements
if (jmxPath != null) {
try {
mbeanServer.unregisterMBean(new ObjectName(mbeanName));
-
} catch (Exception e) {
throw new RuntimeException("Unable to unregister mbean", e);
}
@@ -115,12 +109,12 @@ public class JMXEnabledThreadPoolExecutor extends
ThreadPoolExecutor implements
@Override
public synchronized int getTotalTasks() {
- return totalTasks;
+ return totalTasks.get();
}
@Override
public synchronized double getAverageTaskTime() {
- return (totalTasks == 0) ? 0 : totalTime / totalTasks;
+ return (totalTasks.get() == 0) ? 0 : totalTime.get() /
totalTasks.get();
}
@Override
diff --git
a/server/mailet/mailetcontainer-impl/src/main/java/org/apache/james/mailetcontainer/lib/AbstractStateCompositeProcessor.java
b/server/mailet/mailetcontainer-impl/src/main/java/org/apache/james/mailetcontainer/lib/AbstractStateCompositeProcessor.java
index c36b5d4..a052531 100644
---
a/server/mailet/mailetcontainer-impl/src/main/java/org/apache/james/mailetcontainer/lib/AbstractStateCompositeProcessor.java
+++
b/server/mailet/mailetcontainer-impl/src/main/java/org/apache/james/mailetcontainer/lib/AbstractStateCompositeProcessor.java
@@ -18,12 +18,12 @@
****************************************************************/
package org.apache.james.mailetcontainer.lib;
-import java.util.ArrayList;
-import java.util.Collections;
+import java.util.Collection;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
+import java.util.concurrent.ConcurrentLinkedDeque;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
@@ -49,7 +49,7 @@ import org.slf4j.LoggerFactory;
public abstract class AbstractStateCompositeProcessor implements
MailProcessor, Configurable {
private static final Logger LOGGER =
LoggerFactory.getLogger(AbstractStateCompositeProcessor.class);
- private final List<CompositeProcessorListener> listeners =
Collections.synchronizedList(new ArrayList<>());
+ private final Collection<CompositeProcessorListener> listeners = new
ConcurrentLinkedDeque<>();
private final Map<String, MailProcessor> processors = new HashMap<>();
protected HierarchicalConfiguration<ImmutableNode> config;
diff --git
a/server/mailet/mailetcontainer-impl/src/main/java/org/apache/james/mailetcontainer/lib/AbstractStateMailetProcessor.java
b/server/mailet/mailetcontainer-impl/src/main/java/org/apache/james/mailetcontainer/lib/AbstractStateMailetProcessor.java
index 1f20f45..92fe906 100644
---
a/server/mailet/mailetcontainer-impl/src/main/java/org/apache/james/mailetcontainer/lib/AbstractStateMailetProcessor.java
+++
b/server/mailet/mailetcontainer-impl/src/main/java/org/apache/james/mailetcontainer/lib/AbstractStateMailetProcessor.java
@@ -20,10 +20,10 @@ package org.apache.james.mailetcontainer.lib;
import java.util.ArrayList;
import java.util.Collection;
-import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
+import java.util.concurrent.ConcurrentLinkedDeque;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
@@ -66,7 +66,7 @@ public abstract class AbstractStateMailetProcessor implements
MailProcessor, Con
private MailetContext mailetContext;
private MatcherLoader matcherLoader;
private MailProcessor rootMailProcessor;
- private final List<MailetProcessorListener> listeners =
Collections.synchronizedList(new ArrayList<>());
+ private final Collection<MailetProcessorListener> listeners = new
ConcurrentLinkedDeque<>();
private JMXStateMailetProcessorListener jmxListener;
private boolean enableJmx = true;
private HierarchicalConfiguration<ImmutableNode> config;
@@ -178,7 +178,7 @@ public abstract class AbstractStateMailetProcessor
implements MailProcessor, Con
}
public List<MailetProcessorListener> getListeners() {
- return listeners;
+ return ImmutableList.copyOf(listeners);
}
/**
diff --git
a/server/protocols/protocols-library/src/main/java/org/apache/james/protocols/lib/netty/JMXEnabledOrderedMemoryAwareThreadPoolExecutor.java
b/server/protocols/protocols-library/src/main/java/org/apache/james/protocols/lib/netty/JMXEnabledOrderedMemoryAwareThreadPoolExecutor.java
index d958531..b3cb873 100644
---
a/server/protocols/protocols-library/src/main/java/org/apache/james/protocols/lib/netty/JMXEnabledOrderedMemoryAwareThreadPoolExecutor.java
+++
b/server/protocols/protocols-library/src/main/java/org/apache/james/protocols/lib/netty/JMXEnabledOrderedMemoryAwareThreadPoolExecutor.java
@@ -19,10 +19,10 @@
package org.apache.james.protocols.lib.netty;
import java.lang.management.ManagementFactory;
-import java.util.ArrayList;
-import java.util.Collections;
import java.util.List;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
import javax.management.MBeanServer;
import javax.management.ObjectName;
@@ -36,10 +36,9 @@ import
org.jboss.netty.handler.execution.OrderedMemoryAwareThreadPoolExecutor;
public class JMXEnabledOrderedMemoryAwareThreadPoolExecutor extends
OrderedMemoryAwareThreadPoolExecutor implements
JMXEnabledOrderedMemoryAwareThreadPoolExecutorMBean {
private final String jmxPath;
- private final List<Runnable> inProgress = Collections.synchronizedList(new
ArrayList<>());
private final ThreadLocal<Long> startTime = new ThreadLocal<>();
- private long totalTime;
- private int totalTasks;
+ private final AtomicLong totalTime = new AtomicLong(0);
+ private final AtomicInteger totalTasks = new AtomicInteger(0);
private MBeanServer mbeanServer;
private String mbeanName;
@@ -52,18 +51,14 @@ public class JMXEnabledOrderedMemoryAwareThreadPoolExecutor
extends OrderedMemor
@Override
protected void beforeExecute(Thread t, Runnable r) {
super.beforeExecute(t, r);
- inProgress.add(r);
startTime.set(System.currentTimeMillis());
}
@Override
protected void afterExecute(Runnable r, Throwable t) {
long time = System.currentTimeMillis() - startTime.get();
- synchronized (this) {
- totalTime += time;
- ++totalTasks;
- }
- inProgress.remove(r);
+ totalTime.addAndGet(time);
+ totalTasks.incrementAndGet();
super.afterExecute(r, t);
}
@@ -114,12 +109,12 @@ public class
JMXEnabledOrderedMemoryAwareThreadPoolExecutor extends OrderedMemor
@Override
public synchronized int getTotalTasks() {
- return totalTasks;
+ return totalTasks.get();
}
@Override
public synchronized double getAverageTaskTime() {
- return (totalTasks == 0) ? 0 : totalTime / totalTasks;
+ return (totalTasks.get() == 0) ? 0 : totalTime.get() /
totalTasks.get();
}
@Override
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]