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]

Reply via email to