diff --git a/artemis-cli/src/main/java/org/apache/activemq/artemis/cli/process/ProcessBuilder.java b/artemis-cli/src/main/java/org/apache/activemq/artemis/cli/process/ProcessBuilder.java index 303c956701ef..7eea206276e7 100644 --- a/artemis-cli/src/main/java/org/apache/activemq/artemis/cli/process/ProcessBuilder.java +++ b/artemis-cli/src/main/java/org/apache/activemq/artemis/cli/process/ProcessBuilder.java @@ -21,12 +21,12 @@ import java.io.IOException; import java.io.InputStream; import java.io.InputStreamReader; - -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; public class ProcessBuilder { - static ConcurrentHashSet processes = new ConcurrentHashSet<>(); + static Set processes = ConcurrentHashMap.newKeySet(); static { Runtime.getRuntime().addShutdownHook(new Thread(() -> { diff --git a/artemis-commons/src/main/java/org/apache/activemq/artemis/core/server/NetworkHealthCheck.java b/artemis-commons/src/main/java/org/apache/activemq/artemis/core/server/NetworkHealthCheck.java index d31ef9fcc0d4..25074879e40d 100644 --- a/artemis-commons/src/main/java/org/apache/activemq/artemis/core/server/NetworkHealthCheck.java +++ b/artemis-commons/src/main/java/org/apache/activemq/artemis/core/server/NetworkHealthCheck.java @@ -27,12 +27,12 @@ import java.net.URLConnection; import java.security.PrivilegedAction; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; import org.apache.activemq.artemis.logs.ActiveMQUtilLogger; import org.apache.activemq.artemis.utils.ActiveMQThreadFactory; import org.apache.activemq.artemis.utils.Env; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.activemq.artemis.utils.sm.SecurityManagerShim; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -46,9 +46,9 @@ public class NetworkHealthCheck extends ActiveMQScheduledComponent { private static final Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); - private final Set componentList = new ConcurrentHashSet<>(); - private final Set addresses = new ConcurrentHashSet<>(); - private final Set urls = new ConcurrentHashSet<>(); + private final Set componentList = ConcurrentHashMap.newKeySet(); + private final Set addresses = ConcurrentHashMap.newKeySet(); + private final Set urls = ConcurrentHashMap.newKeySet(); private NetworkInterface networkInterface; public static final String IPV6_DEFAULT_COMMAND = Env.isWindowsOs() ? "ping -n 1 -w %d000 %s" : "ping6 -c 1 %2$s"; diff --git a/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/collections/ConcurrentHashSet.java b/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/collections/ConcurrentHashSet.java index 93b0b809dd2a..663225199c2c 100644 --- a/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/collections/ConcurrentHashSet.java +++ b/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/collections/ConcurrentHashSet.java @@ -26,6 +26,7 @@ *

* Offers same concurrency as ConcurrentHashMap but for a Set */ +@Deprecated(forRemoval = true) public class ConcurrentHashSet extends AbstractSet implements ConcurrentSet { private final ConcurrentMap theMap; diff --git a/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/collections/ConcurrentSet.java b/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/collections/ConcurrentSet.java index bab7aa37477f..3ff6f6f74a88 100644 --- a/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/collections/ConcurrentSet.java +++ b/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/collections/ConcurrentSet.java @@ -23,6 +23,7 @@ * * @param The generic class */ +@Deprecated(forRemoval = true) public interface ConcurrentSet extends Set { boolean addIfAbsent(E o); diff --git a/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/critical/CriticalAnalyzerImpl.java b/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/critical/CriticalAnalyzerImpl.java index a825fe894980..d6563f8f5c26 100644 --- a/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/critical/CriticalAnalyzerImpl.java +++ b/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/critical/CriticalAnalyzerImpl.java @@ -19,12 +19,13 @@ import java.security.PrivilegedAction; import java.util.ConcurrentModificationException; import java.util.List; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.TimeUnit; import org.apache.activemq.artemis.core.server.ActiveMQScheduledComponent; import org.apache.activemq.artemis.utils.ActiveMQThreadFactory; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.activemq.artemis.utils.sm.SecurityManagerShim; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -79,7 +80,7 @@ public void clear() { private List actions = new CopyOnWriteArrayList<>(); - private final ConcurrentHashSet components = new ConcurrentHashSet<>(); + private final Set components = ConcurrentHashMap.newKeySet(); @Override public int getNumberOfComponents() { diff --git a/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/uri/FluentPropertyBeanIntrospectorWithIgnores.java b/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/uri/FluentPropertyBeanIntrospectorWithIgnores.java index 31ff14d37891..7729a48e92f5 100644 --- a/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/uri/FluentPropertyBeanIntrospectorWithIgnores.java +++ b/artemis-commons/src/main/java/org/apache/activemq/artemis/utils/uri/FluentPropertyBeanIntrospectorWithIgnores.java @@ -21,9 +21,10 @@ import java.beans.PropertyDescriptor; import java.lang.reflect.Method; import java.util.Locale; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import org.apache.activemq.artemis.api.core.Pair; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.commons.beanutils.FluentPropertyBeanIntrospector; import org.apache.commons.beanutils.IntrospectionContext; import org.slf4j.Logger; @@ -34,7 +35,7 @@ public class FluentPropertyBeanIntrospectorWithIgnores extends FluentPropertyBea static Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); - private static ConcurrentHashSet> ignores = new ConcurrentHashSet<>(); + private static Set> ignores = ConcurrentHashMap.newKeySet(); public static void addIgnore(String className, String methodName) { logger.trace("Adding ignore on {}/{}", className, methodName); diff --git a/artemis-core-client/src/main/java/org/apache/activemq/artemis/core/client/impl/ClientSessionFactoryImpl.java b/artemis-core-client/src/main/java/org/apache/activemq/artemis/core/client/impl/ClientSessionFactoryImpl.java index 999fe9045654..967fac057ece 100644 --- a/artemis-core-client/src/main/java/org/apache/activemq/artemis/core/client/impl/ClientSessionFactoryImpl.java +++ b/artemis-core-client/src/main/java/org/apache/activemq/artemis/core/client/impl/ClientSessionFactoryImpl.java @@ -20,10 +20,10 @@ import java.lang.ref.WeakReference; import java.security.PrivilegedAction; import java.util.ArrayList; -import java.util.Collections; import java.util.HashSet; import java.util.List; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.Executor; import java.util.concurrent.Future; @@ -73,7 +73,6 @@ import org.apache.activemq.artemis.utils.UUIDGenerator; import org.apache.activemq.artemis.utils.actors.ArtemisExecutor; import org.apache.activemq.artemis.utils.actors.OrderedExecutorFactory; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.activemq.artemis.utils.sm.SecurityManagerShim; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -104,7 +103,7 @@ public class ClientSessionFactoryImpl implements ClientSessionFactoryInternal, C private final long connectionTTL; - private final Set sessions = new ConcurrentHashSet<>(); + private final Set sessions = ConcurrentHashMap.newKeySet(); private final Object createSessionLock = new Object(); @@ -138,9 +137,9 @@ public class ClientSessionFactoryImpl implements ClientSessionFactoryInternal, C private int failoverAttempts; - private final Set listeners = new ConcurrentHashSet<>(); + private final Set listeners = ConcurrentHashMap.newKeySet(); - private final Set failoverListeners = new ConcurrentHashSet<>(); + private final Set failoverListeners = ConcurrentHashMap.newKeySet(); private Connector connector; @@ -155,7 +154,7 @@ public class ClientSessionFactoryImpl implements ClientSessionFactoryInternal, C private volatile boolean closed; - public static final Set CLOSE_RUNNABLES = Collections.synchronizedSet(new HashSet<>()); + public static final Set CLOSE_RUNNABLES = ConcurrentHashMap.newKeySet(); private final ConfirmationWindowWarning confirmationWindowWarning; diff --git a/artemis-core-client/src/test/java/org/apache/activemq/artemis/util/TimeAndCounterIDGeneratorTest.java b/artemis-core-client/src/test/java/org/apache/activemq/artemis/util/TimeAndCounterIDGeneratorTest.java index ee1ee2833522..2925d76e7697 100644 --- a/artemis-core-client/src/test/java/org/apache/activemq/artemis/util/TimeAndCounterIDGeneratorTest.java +++ b/artemis-core-client/src/test/java/org/apache/activemq/artemis/util/TimeAndCounterIDGeneratorTest.java @@ -20,11 +20,12 @@ import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.fail; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import org.apache.activemq.artemis.utils.TimeAndCounterIDGenerator; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.junit.jupiter.api.Test; public class TimeAndCounterIDGeneratorTest { @@ -66,7 +67,7 @@ public void testCalculationRefresh() { @Test public void testCalculationOnMultiThread() throws Throwable { - final ConcurrentHashSet hashSet = new ConcurrentHashSet<>(); + final Set hashSet = ConcurrentHashMap.newKeySet(); final TimeAndCounterIDGenerator seq = new TimeAndCounterIDGenerator(); diff --git a/artemis-jms-client/src/main/java/org/apache/activemq/artemis/jms/client/ActiveMQConnection.java b/artemis-jms-client/src/main/java/org/apache/activemq/artemis/jms/client/ActiveMQConnection.java index dc0e2b9aed8c..ec835de28e4b 100644 --- a/artemis-jms-client/src/main/java/org/apache/activemq/artemis/jms/client/ActiveMQConnection.java +++ b/artemis-jms-client/src/main/java/org/apache/activemq/artemis/jms/client/ActiveMQConnection.java @@ -20,6 +20,7 @@ import java.security.PrivilegedAction; import java.util.HashSet; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -55,7 +56,6 @@ import org.apache.activemq.artemis.utils.ActiveMQThreadFactory; import org.apache.activemq.artemis.utils.UUIDGenerator; import org.apache.activemq.artemis.utils.VersionLoader; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.activemq.artemis.utils.sm.SecurityManagerShim; /** @@ -82,9 +82,9 @@ public class ActiveMQConnection extends ActiveMQConnectionForContextImpl impleme private final int connectionType; - private final Set sessions = new ConcurrentHashSet<>(); + private final Set sessions = ConcurrentHashMap.newKeySet(); - private final Set tempQueues = new ConcurrentHashSet<>(); + private final Set tempQueues = ConcurrentHashMap.newKeySet(); private volatile boolean hasNoLocal; diff --git a/artemis-jms-client/src/main/java/org/apache/activemq/artemis/jms/client/ThreadAwareContext.java b/artemis-jms-client/src/main/java/org/apache/activemq/artemis/jms/client/ThreadAwareContext.java index de15fcbd0c90..d54b44ebfd91 100644 --- a/artemis-jms-client/src/main/java/org/apache/activemq/artemis/jms/client/ThreadAwareContext.java +++ b/artemis-jms-client/src/main/java/org/apache/activemq/artemis/jms/client/ThreadAwareContext.java @@ -18,8 +18,7 @@ import javax.jms.IllegalStateException; import java.util.Set; - -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; +import java.util.concurrent.ConcurrentHashMap; /** * Restricts what can be called on context passed in wrapped CompletionListener. @@ -39,7 +38,7 @@ public class ThreadAwareContext { * Use a set because JMSContext can create more than one JMSConsumer to receive asynchronously from different * destinations. */ - private final Set messageListenerThreads = new ConcurrentHashSet<>(); + private final Set messageListenerThreads = ConcurrentHashMap.newKeySet(); /** * Sets current thread to the context diff --git a/artemis-journal/src/main/java/org/apache/activemq/artemis/core/journal/impl/JournalImpl.java b/artemis-journal/src/main/java/org/apache/activemq/artemis/core/journal/impl/JournalImpl.java index 737638c7f054..775c687b32e6 100644 --- a/artemis-journal/src/main/java/org/apache/activemq/artemis/core/journal/impl/JournalImpl.java +++ b/artemis-journal/src/main/java/org/apache/activemq/artemis/core/journal/impl/JournalImpl.java @@ -34,6 +34,8 @@ import java.util.LinkedList; import java.util.List; import java.util.Map; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.Executor; import java.util.concurrent.RejectedExecutionException; @@ -87,7 +89,6 @@ import org.apache.activemq.artemis.utils.SimpleFuture; import org.apache.activemq.artemis.utils.SimpleFutureImpl; import org.apache.activemq.artemis.utils.actors.OrderedExecutorFactory; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.activemq.artemis.utils.collections.ConcurrentLongHashMap; import org.apache.activemq.artemis.utils.collections.LongHashSet; import org.apache.activemq.artemis.utils.collections.SparseArrayLinkedList; @@ -313,7 +314,7 @@ public JournalImpl setHistoryFolder(File historyFolder, long maxBytes, long peri private Executor appendExecutor = null; - private final ConcurrentHashSet latches = new ConcurrentHashSet<>(); + private final Set latches = ConcurrentHashMap.newKeySet(); private final ExecutorFactory providedIOThreadPool; protected ExecutorFactory ioExecutorFactory; diff --git a/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/connect/bridge/AMQPBridgeManagers.java b/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/connect/bridge/AMQPBridgeManagers.java index 258b20f43b13..c243ee3fbff6 100644 --- a/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/connect/bridge/AMQPBridgeManagers.java +++ b/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/connect/bridge/AMQPBridgeManagers.java @@ -18,11 +18,11 @@ import java.lang.invoke.MethodHandles; import java.util.Collection; +import java.util.concurrent.ConcurrentHashMap; import org.apache.activemq.artemis.api.core.ActiveMQException; import org.apache.activemq.artemis.core.config.amqpBrokerConnectivity.AMQPBridgeBrokerConnectionElement; import org.apache.activemq.artemis.protocol.amqp.connect.AMQPBrokerConnection; import org.apache.activemq.artemis.protocol.amqp.proton.AMQPSessionContext; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -33,7 +33,7 @@ public class AMQPBridgeManagers { private static final Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); - private final Collection bridgeManagers = new ConcurrentHashSet<>(); + private final Collection bridgeManagers = ConcurrentHashMap.newKeySet(); private final AMQPBrokerConnection brokerConnection; public AMQPBridgeManagers(AMQPBrokerConnection brokerConnection) { diff --git a/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/connect/federation/AMQPFederationAddressBindingsConsumer.java b/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/connect/federation/AMQPFederationAddressBindingsConsumer.java index 7abe3021c63d..681de6f1ce8f 100644 --- a/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/connect/federation/AMQPFederationAddressBindingsConsumer.java +++ b/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/connect/federation/AMQPFederationAddressBindingsConsumer.java @@ -21,6 +21,7 @@ import java.util.ArrayList; import java.util.Collection; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import org.apache.activemq.artemis.api.core.Message; import org.apache.activemq.artemis.core.io.IOCallback; @@ -34,7 +35,6 @@ import org.apache.activemq.artemis.protocol.amqp.connect.federation.AMQPFederationMetrics.ConsumerMetrics; import org.apache.activemq.artemis.protocol.amqp.federation.FederationConsumerInfo; import org.apache.activemq.artemis.protocol.amqp.proton.AMQPSessionContext; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.qpid.proton.amqp.messaging.Accepted; import org.apache.qpid.proton.engine.Delivery; import org.apache.qpid.proton.engine.Receiver; @@ -59,7 +59,7 @@ public final class AMQPFederationAddressBindingsConsumer extends AMQPFederationA private static final Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); - private final Set bindings = new ConcurrentHashSet<>(); + private final Set bindings = ConcurrentHashMap.newKeySet(); private final PostOffice postOffice; private final StorageManager storageManager; diff --git a/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/proton/AMQPConnectionContext.java b/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/proton/AMQPConnectionContext.java index e785bb4ed83c..fdd1716028b7 100644 --- a/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/proton/AMQPConnectionContext.java +++ b/artemis-protocols/artemis-amqp-protocol/src/main/java/org/apache/activemq/artemis/protocol/amqp/proton/AMQPConnectionContext.java @@ -64,7 +64,6 @@ import org.apache.activemq.artemis.spi.core.remoting.ReadyListener; import org.apache.activemq.artemis.utils.ByteUtil; import org.apache.activemq.artemis.utils.VersionLoader; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.qpid.proton.amqp.Symbol; import org.apache.qpid.proton.amqp.messaging.Source; import org.apache.qpid.proton.amqp.messaging.TerminusExpiryPolicy; @@ -127,7 +126,7 @@ public void enableAutoRead() { private final Symbol[] desiredCapabilities; private final ScheduledExecutorService scheduledPool; private final Map linkCloseListeners = new ConcurrentHashMap<>(); - private final Set> remoteOpenedListeners = new ConcurrentHashSet<>(); + private final Set> remoteOpenedListeners = ConcurrentHashMap.newKeySet(); private final Map sessions = new ConcurrentHashMap<>(); diff --git a/artemis-protocols/artemis-openwire-protocol/src/main/java/org/apache/activemq/artemis/core/protocol/openwire/OpenWireConnection.java b/artemis-protocols/artemis-openwire-protocol/src/main/java/org/apache/activemq/artemis/core/protocol/openwire/OpenWireConnection.java index 9473d0235109..3ea9e3b42241 100644 --- a/artemis-protocols/artemis-openwire-protocol/src/main/java/org/apache/activemq/artemis/core/protocol/openwire/OpenWireConnection.java +++ b/artemis-protocols/artemis-openwire-protocol/src/main/java/org/apache/activemq/artemis/core/protocol/openwire/OpenWireConnection.java @@ -83,7 +83,6 @@ import org.apache.activemq.artemis.spi.core.remoting.Connection; import org.apache.activemq.artemis.utils.UUIDGenerator; import org.apache.activemq.artemis.utils.actors.ThresholdActor; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.activemq.command.ActiveMQDestination; import org.apache.activemq.command.ActiveMQMessage; import org.apache.activemq.command.ActiveMQTempQueue; @@ -208,7 +207,7 @@ public class OpenWireConnection extends AbstractRemotingConnection implements Se private long maxInactivityDuration; private volatile ThresholdActor openWireActor; - private final Set knownDestinations = new ConcurrentHashSet<>(); + private final Set knownDestinations = ConcurrentHashMap.newKeySet(); private final AtomicBoolean disableTtl = new AtomicBoolean(false); diff --git a/artemis-ra/src/main/java/org/apache/activemq/artemis/ra/ActiveMQRAConnectionManager.java b/artemis-ra/src/main/java/org/apache/activemq/artemis/ra/ActiveMQRAConnectionManager.java index 9519ab979108..f826cb5730dd 100644 --- a/artemis-ra/src/main/java/org/apache/activemq/artemis/ra/ActiveMQRAConnectionManager.java +++ b/artemis-ra/src/main/java/org/apache/activemq/artemis/ra/ActiveMQRAConnectionManager.java @@ -22,13 +22,14 @@ import javax.resource.spi.ManagedConnection; import javax.resource.spi.ManagedConnectionFactory; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.IOException; import java.io.ObjectInputStream; import java.lang.invoke.MethodHandles; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; /** * The connection manager used in non-managed environments. @@ -42,7 +43,7 @@ public ActiveMQRAConnectionManager() { logger.trace("constructor()"); } - transient ConcurrentHashSet connections = new ConcurrentHashSet<>(); + transient Set connections = ConcurrentHashMap.newKeySet(); /** * Allocates a connection @@ -81,6 +82,6 @@ public void stop() { */ private void readObject(ObjectInputStream in) throws IOException, ClassNotFoundException { in.defaultReadObject(); - connections = new ConcurrentHashSet<>(); + connections = ConcurrentHashMap.newKeySet(); } } diff --git a/artemis-ra/src/main/java/org/apache/activemq/artemis/ra/ActiveMQRAManagedConnection.java b/artemis-ra/src/main/java/org/apache/activemq/artemis/ra/ActiveMQRAManagedConnection.java index 761525aa71d4..4de691db71fe 100644 --- a/artemis-ra/src/main/java/org/apache/activemq/artemis/ra/ActiveMQRAManagedConnection.java +++ b/artemis-ra/src/main/java/org/apache/activemq/artemis/ra/ActiveMQRAManagedConnection.java @@ -39,10 +39,10 @@ import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; -import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.locks.ReentrantLock; @@ -73,9 +73,9 @@ public final class ActiveMQRAManagedConnection implements ManagedConnection, Exc private final AtomicBoolean isDestroyed = new AtomicBoolean(false); - private final List eventListeners; + private final List eventListeners = Collections.synchronizedList(new ArrayList<>()); - private final Set handles; + private final Set handles = ConcurrentHashMap.newKeySet(); private ReentrantLock lock = new ReentrantLock(); @@ -111,8 +111,6 @@ public ActiveMQRAManagedConnection(final ActiveMQRAManagedConnectionFactory mcf, this.ra = ra; this.userName = userName; this.password = password; - eventListeners = Collections.synchronizedList(new ArrayList<>()); - handles = Collections.synchronizedSet(new HashSet<>()); connection = null; nonXAsession = null; diff --git a/artemis-ra/src/main/java/org/apache/activemq/artemis/ra/recovery/RecoveryManager.java b/artemis-ra/src/main/java/org/apache/activemq/artemis/ra/recovery/RecoveryManager.java index 8bb19f9d08c2..a77a89ecd223 100644 --- a/artemis-ra/src/main/java/org/apache/activemq/artemis/ra/recovery/RecoveryManager.java +++ b/artemis-ra/src/main/java/org/apache/activemq/artemis/ra/recovery/RecoveryManager.java @@ -23,12 +23,12 @@ import java.util.Map; import java.util.ServiceLoader; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import org.apache.activemq.artemis.jms.client.ActiveMQConnectionFactory; import org.apache.activemq.artemis.service.extensions.xa.recovery.ActiveMQRegistry; import org.apache.activemq.artemis.service.extensions.xa.recovery.ActiveMQRegistryImpl; import org.apache.activemq.artemis.service.extensions.xa.recovery.XARecoveryConfig; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -42,7 +42,7 @@ public final class RecoveryManager implements Serializable { private static final String RESOURCE_RECOVERY_CLASS_NAMES = "org.jboss.as.messaging.jms.AS7RecoveryRegistry;" + "org.jboss.as.integration.activemq.recovery.AS5RecoveryRegistry"; - private transient Set resources = new ConcurrentHashSet<>(); + private transient Set resources = ConcurrentHashMap.newKeySet(); public void start(final boolean useAutoRecovery) { if (useAutoRecovery) { @@ -116,6 +116,6 @@ public Set getResources() { */ private void readObject(ObjectInputStream in) throws IOException, ClassNotFoundException { in.defaultReadObject(); - resources = new ConcurrentHashSet<>(); + resources = ConcurrentHashMap.newKeySet(); } } diff --git a/artemis-server/src/main/java/org/apache/activemq/artemis/core/paging/impl/PagingManagerImpl.java b/artemis-server/src/main/java/org/apache/activemq/artemis/core/paging/impl/PagingManagerImpl.java index 4ce1cb3b3fcf..ffe64841a3c4 100644 --- a/artemis-server/src/main/java/org/apache/activemq/artemis/core/paging/impl/PagingManagerImpl.java +++ b/artemis-server/src/main/java/org/apache/activemq/artemis/core/paging/impl/PagingManagerImpl.java @@ -51,7 +51,6 @@ import org.apache.activemq.artemis.utils.ByteUtil; import org.apache.activemq.artemis.utils.CompositeAddress; import org.apache.activemq.artemis.utils.SizeAwareMetric; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.activemq.artemis.utils.runnables.AtomicRunnable; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -76,7 +75,7 @@ public final class PagingManagerImpl implements PagingManager { */ private final ReentrantReadWriteLock syncLock = new ReentrantReadWriteLock(); - private final Set blockedStored = new ConcurrentHashSet<>(); + private final Set blockedStored = ConcurrentHashMap.newKeySet(); private final ConcurrentMap stores = new ConcurrentHashMap<>(); diff --git a/artemis-server/src/main/java/org/apache/activemq/artemis/core/security/impl/SecurityStoreImpl.java b/artemis-server/src/main/java/org/apache/activemq/artemis/core/security/impl/SecurityStoreImpl.java index 3370faca3c2c..090ba5776ed1 100644 --- a/artemis-server/src/main/java/org/apache/activemq/artemis/core/security/impl/SecurityStoreImpl.java +++ b/artemis-server/src/main/java/org/apache/activemq/artemis/core/security/impl/SecurityStoreImpl.java @@ -26,6 +26,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLongFieldUpdater; import java.util.stream.Collectors; +import java.util.concurrent.ConcurrentHashMap; import com.github.benmanes.caffeine.cache.Cache; import com.github.benmanes.caffeine.cache.Caffeine; @@ -58,7 +59,6 @@ import org.apache.activemq.artemis.utils.ByteUtil; import org.apache.activemq.artemis.utils.CertificateUtil; import org.apache.activemq.artemis.utils.CompositeAddress; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.activemq.artemis.utils.collections.TypedProperties; import org.apache.activemq.artemis.utils.sm.SecurityManagerShim; import org.slf4j.Logger; @@ -81,7 +81,7 @@ public class SecurityStoreImpl implements SecurityStore, HierarchicalRepositoryC private final ActiveMQSecurityManager securityManager; - private final Cache> authorizationCache; + private final Cache> authorizationCache; private final Cache> authenticationCache; @@ -357,13 +357,13 @@ public boolean hasPermission(final SimpleString address, if (validated && user != null) { // if we get here we're granted, add to the cache - ConcurrentHashSet set; + final Set set; String key = createAuthorizationCacheKey(user, checkType); - ConcurrentHashSet act = getAuthorizationCacheEntry(key); + Set act = getAuthorizationCacheEntry(key); if (act != null) { set = act; } else { - set = new ConcurrentHashSet<>(); + set = ConcurrentHashMap.newKeySet(); putAuthorizationCacheEntry(set, key); } set.add(Objects.requireNonNullElse(fqqn, bareAddress)); @@ -566,18 +566,18 @@ private Pair getAuthenticationCacheEntry(String key) { } } - private void putAuthorizationCacheEntry(ConcurrentHashSet value, String key) { + private void putAuthorizationCacheEntry(Set value, String key) { if (authorizationCache != null) { authorizationCache.put(key, value); logger.trace("Put into authz cache; key: {}; value: {}", key, value); } } - private ConcurrentHashSet getAuthorizationCacheEntry(String key) { + private Set getAuthorizationCacheEntry(String key) { if (authorizationCache == null) { return null; } else { - ConcurrentHashSet value = authorizationCache.getIfPresent(key); + Set value = authorizationCache.getIfPresent(key); logger.trace("Get from authz cache; key: {}; value: {}", key, value); return value; } @@ -616,7 +616,7 @@ public long getAuthorizationCacheSize() { private boolean checkAuthorizationCache(final SimpleString dest, final String user, final CheckType checkType) { boolean granted = false; - ConcurrentHashSet act = getAuthorizationCacheEntry(createAuthorizationCacheKey(user, checkType)); + Set act = getAuthorizationCacheEntry(createAuthorizationCacheKey(user, checkType)); if (act != null) { granted = act.contains(dest); } @@ -671,7 +671,7 @@ public Cache> getAuthenticationCache() { return authenticationCache; } - public Cache> getAuthorizationCache() { + public Cache> getAuthorizationCache() { return authorizationCache; } diff --git a/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/cluster/ClusterManager.java b/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/cluster/ClusterManager.java index ecac84215bfb..91dfd672ae87 100644 --- a/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/cluster/ClusterManager.java +++ b/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/cluster/ClusterManager.java @@ -27,6 +27,7 @@ import java.util.Map; import java.util.Objects; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.Executor; import java.util.concurrent.ScheduledExecutorService; @@ -65,7 +66,6 @@ import org.apache.activemq.artemis.spi.core.remoting.Acceptor; import org.apache.activemq.artemis.utils.ExecutorFactory; import org.apache.activemq.artemis.utils.FutureLatch; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.lang.invoke.MethodHandles; @@ -145,7 +145,7 @@ enum State { // the cluster connections which links this node to other cluster nodes private final Map clusterConnections = new HashMap<>(); - private final Set clusterLocators = new ConcurrentHashSet<>(); + private final Set clusterLocators = ConcurrentHashMap.newKeySet(); private final Executor executor; diff --git a/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/group/impl/RemoteGroupingHandler.java b/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/group/impl/RemoteGroupingHandler.java index 99216c4d78ba..b7d2d0c27b3a 100644 --- a/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/group/impl/RemoteGroupingHandler.java +++ b/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/group/impl/RemoteGroupingHandler.java @@ -19,6 +19,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.TimeUnit; @@ -35,7 +36,6 @@ import org.apache.activemq.artemis.core.server.management.ManagementService; import org.apache.activemq.artemis.core.server.management.Notification; import org.apache.activemq.artemis.utils.ExecutorFactory; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.activemq.artemis.utils.collections.TypedProperties; /** @@ -60,7 +60,7 @@ public final class RemoteGroupingHandler extends GroupHandlingAbstract { private final ConcurrentMap> groupMap = new ConcurrentHashMap<>(); - private final ConcurrentHashSet pendingNotifications = new ConcurrentHashSet<>(); + private final Set pendingNotifications = ConcurrentHashMap.newKeySet(); private boolean started = false; diff --git a/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/impl/ActiveMQServerImpl.java b/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/impl/ActiveMQServerImpl.java index 6d95de2fb8c7..0154a605accb 100644 --- a/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/impl/ActiveMQServerImpl.java +++ b/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/impl/ActiveMQServerImpl.java @@ -217,7 +217,6 @@ import org.apache.activemq.artemis.utils.UUID; import org.apache.activemq.artemis.utils.VersionLoader; import org.apache.activemq.artemis.utils.actors.OrderedExecutorFactory; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.activemq.artemis.utils.critical.CriticalAction; import org.apache.activemq.artemis.utils.critical.CriticalAnalyzer; import org.apache.activemq.artemis.utils.critical.CriticalAnalyzerImpl; @@ -365,15 +364,15 @@ public class ActiveMQServerImpl implements ActiveMQServer { private final Map brokerConnectionMap = new ConcurrentHashMap<>(); - private final Set activateCallbacks = new ConcurrentHashSet<>(); + private final Set activateCallbacks = ConcurrentHashMap.newKeySet(); - private final Set activationFailureListeners = new ConcurrentHashSet<>(); + private final Set activationFailureListeners = ConcurrentHashMap.newKeySet(); - private final Set ioCriticalErrorListeners = new ConcurrentHashSet<>(); + private final Set ioCriticalErrorListeners = ConcurrentHashMap.newKeySet(); - private final Set postQueueCreationCallbacks = new ConcurrentHashSet<>(); + private final Set postQueueCreationCallbacks = ConcurrentHashMap.newKeySet(); - private final Set postQueueDeletionCallbacks = new ConcurrentHashSet<>(); + private final Set postQueueDeletionCallbacks = ConcurrentHashMap.newKeySet(); private volatile GroupingHandler groupingHandler; diff --git a/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/impl/QueueImpl.java b/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/impl/QueueImpl.java index 0e867595e2b2..5b18c3a9b5de 100644 --- a/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/impl/QueueImpl.java +++ b/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/impl/QueueImpl.java @@ -32,6 +32,7 @@ import java.util.Objects; import java.util.Random; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; @@ -117,7 +118,6 @@ import org.apache.activemq.artemis.utils.ReusableLatch; import org.apache.activemq.artemis.utils.SizeAwareMetric; import org.apache.activemq.artemis.utils.actors.ArtemisExecutor; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.activemq.artemis.utils.collections.LinkedListIterator; import org.apache.activemq.artemis.utils.collections.NodeStoreFactory; import org.apache.activemq.artemis.utils.collections.PriorityLinkedList; @@ -342,7 +342,7 @@ public void setSwept(boolean swept) { */ private final Object directDeliveryGuard = new Object(); - private final ConcurrentHashSet lingerSessionIds = new ConcurrentHashSet<>(); + private final Set lingerSessionIds = ConcurrentHashMap.newKeySet(); public String debug() { StringWriter str = new StringWriter(); diff --git a/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/management/impl/ManagementServiceImpl.java b/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/management/impl/ManagementServiceImpl.java index af6222480ffc..65981610b322 100644 --- a/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/management/impl/ManagementServiceImpl.java +++ b/artemis-server/src/main/java/org/apache/activemq/artemis/core/server/management/impl/ManagementServiceImpl.java @@ -30,6 +30,7 @@ import java.util.List; import java.util.Objects; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ScheduledExecutorService; import java.util.function.Predicate; import java.util.regex.Pattern; @@ -115,7 +116,6 @@ import org.apache.activemq.artemis.core.settings.impl.AddressSettings; import org.apache.activemq.artemis.core.transaction.ResourceManager; import org.apache.activemq.artemis.spi.core.remoting.Acceptor; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.activemq.artemis.utils.collections.TypedProperties; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -162,7 +162,7 @@ public class ManagementServiceImpl implements ManagementService { private boolean notificationsEnabled; - private final Set listeners = new ConcurrentHashSet<>(); + private final Set listeners = ConcurrentHashMap.newKeySet(); private final ObjectNameBuilder objectNameBuilder; @@ -170,7 +170,7 @@ public class ManagementServiceImpl implements ManagementService { private final Pattern viewPermissionMatcher; - private final Set registeredNames = new ConcurrentHashSet<>(); + private final Set registeredNames = ConcurrentHashMap.newKeySet(); public ManagementServiceImpl(final MBeanServer mbeanServer, final Configuration configuration) { this.mbeanServer = mbeanServer; diff --git a/artemis-server/src/main/java/org/apache/activemq/artemis/core/settings/impl/HierarchicalObjectRepository.java b/artemis-server/src/main/java/org/apache/activemq/artemis/core/settings/impl/HierarchicalObjectRepository.java index be6737f27289..0ca275983ee8 100644 --- a/artemis-server/src/main/java/org/apache/activemq/artemis/core/settings/impl/HierarchicalObjectRepository.java +++ b/artemis-server/src/main/java/org/apache/activemq/artemis/core/settings/impl/HierarchicalObjectRepository.java @@ -38,7 +38,6 @@ import org.apache.activemq.artemis.core.settings.HierarchicalRepositoryChangeListener; import org.apache.activemq.artemis.core.settings.Mergeable; import org.apache.activemq.artemis.utils.CompositeAddress; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -106,7 +105,7 @@ public class HierarchicalObjectRepository implements HierarchicalRepository listeners = new ConcurrentHashSet<>(); + private final Set listeners = ConcurrentHashMap.newKeySet(); public HierarchicalObjectRepository() { this(null); diff --git a/artemis-server/src/test/java/org/apache/activemq/artemis/core/server/group/impl/ClusteredResetMockTest.java b/artemis-server/src/test/java/org/apache/activemq/artemis/core/server/group/impl/ClusteredResetMockTest.java index 2858ddb278be..953524bf89da 100644 --- a/artemis-server/src/test/java/org/apache/activemq/artemis/core/server/group/impl/ClusteredResetMockTest.java +++ b/artemis-server/src/test/java/org/apache/activemq/artemis/core/server/group/impl/ClusteredResetMockTest.java @@ -21,6 +21,7 @@ import java.util.HashSet; import java.util.List; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.function.Predicate; @@ -79,7 +80,6 @@ import org.apache.activemq.artemis.spi.core.remoting.Acceptor; import org.apache.activemq.artemis.tests.util.ServerTestBase; import org.apache.activemq.artemis.utils.ReusableLatch; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.junit.jupiter.api.Test; /** @@ -190,7 +190,7 @@ public void run() { class FakeManagement implements ManagementService { - public ConcurrentHashSet pendingNotifications = new ConcurrentHashSet<>(); + public Set pendingNotifications = ConcurrentHashMap.newKeySet(); final ReusableLatch latch; diff --git a/tests/artemis-test-support/src/main/java/org/apache/activemq/artemis/tests/integration/stomp/util/AbstractStompClientConnection.java b/tests/artemis-test-support/src/main/java/org/apache/activemq/artemis/tests/integration/stomp/util/AbstractStompClientConnection.java index cf6d129205eb..21f859b700d6 100644 --- a/tests/artemis-test-support/src/main/java/org/apache/activemq/artemis/tests/integration/stomp/util/AbstractStompClientConnection.java +++ b/tests/artemis-test-support/src/main/java/org/apache/activemq/artemis/tests/integration/stomp/util/AbstractStompClientConnection.java @@ -23,7 +23,9 @@ import java.util.ArrayList; import java.util.List; import java.util.Objects; +import java.util.Set; import java.util.concurrent.BlockingQueue; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; @@ -31,7 +33,6 @@ import io.netty.buffer.Unpooled; import io.netty.channel.ChannelFuture; import org.apache.activemq.artemis.core.protocol.stomp.Stomp; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.activemq.transport.netty.NettyTransport; import org.apache.activemq.transport.netty.NettyTransportFactory; import org.apache.activemq.transport.netty.NettyTransportListener; @@ -54,7 +55,7 @@ public abstract class AbstractStompClientConnection implements StompClientConnec //protected ReaderThread readerThread; protected String scheme; - private static final ConcurrentHashSet connections = new ConcurrentHashSet<>(); + private static final Set connections = ConcurrentHashMap.newKeySet(); public static final void tearDownConnections() { diff --git a/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/client/ConsumerTest.java b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/client/ConsumerTest.java index 6bfceb33b607..e2d6c3d118de 100644 --- a/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/client/ConsumerTest.java +++ b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/client/ConsumerTest.java @@ -36,6 +36,7 @@ import java.util.Collection; import java.util.Set; import java.util.concurrent.BrokenBarrierException; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.CyclicBarrier; import java.util.concurrent.TimeUnit; @@ -71,7 +72,6 @@ import org.apache.activemq.artemis.tests.util.ActiveMQTestBase; import org.apache.activemq.artemis.tests.util.Wait; import org.apache.activemq.artemis.utils.ByteUtil; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.apache.qpid.jms.JmsConnectionFactory; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.TestTemplate; @@ -913,7 +913,7 @@ public void testNoReceiveWithListener() throws Exception { @TestTemplate public void testReceiveAndResend() throws Exception { - final Set sessions = new ConcurrentHashSet<>(); + final Set sessions = ConcurrentHashMap.newKeySet(); final AtomicInteger errors = new AtomicInteger(0); final SimpleString QUEUE_RESPONSE = SimpleString.of("QUEUE_RESPONSE"); diff --git a/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/client/SlowConsumerTest.java b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/client/SlowConsumerTest.java index 9b934ea074f4..a0d2d5b8a273 100644 --- a/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/client/SlowConsumerTest.java +++ b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/client/SlowConsumerTest.java @@ -18,6 +18,7 @@ import java.lang.invoke.MethodHandles; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; @@ -46,7 +47,6 @@ import org.apache.activemq.artemis.tests.util.ActiveMQTestBase; import org.apache.activemq.artemis.tests.util.Wait; import org.apache.activemq.artemis.utils.RandomUtil; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.slf4j.Logger; @@ -454,7 +454,7 @@ private void testMinuteKilled(boolean netty) throws Exception { session.commit(); - ConcurrentHashSet receivedMessages = new ConcurrentHashSet<>(); + Set receivedMessages = ConcurrentHashMap.newKeySet(); FixedRateConsumer consumer = new FixedRateConsumer(40, MESSAGES_PER_MINUTE, receivedMessages, sf, QUEUE, 0); consumer.start(); @@ -496,7 +496,7 @@ public void testDaysKilled() throws Exception { session.commit(); - ConcurrentHashSet receivedMessages = new ConcurrentHashSet<>(); + Set receivedMessages = ConcurrentHashMap.newKeySet(); FixedRateConsumer consumer = new FixedRateConsumer(30, MESSAGES_PER_MINUTE, receivedMessages, sf, QUEUE, 0); consumer.start(); @@ -548,7 +548,7 @@ public void testDaysKilledPaging() throws Exception { session.commit(); - ConcurrentHashSet receivedMessages = new ConcurrentHashSet<>(); + Set receivedMessages = ConcurrentHashMap.newKeySet(); FixedRateConsumer consumer = new FixedRateConsumer(30, MESSAGES_PER_MINUTE, receivedMessages, sf, QUEUE, 0); consumer.start(); @@ -589,7 +589,7 @@ public void testDaysSurviving() throws Exception { session.commit(); - ConcurrentHashSet receivedMessages = new ConcurrentHashSet<>(); + Set receivedMessages = ConcurrentHashMap.newKeySet(); FixedRateConsumer consumer = new FixedRateConsumer(70, MESSAGES_PER_MINUTE, receivedMessages, sf, QUEUE, 0); consumer.start(); @@ -634,7 +634,7 @@ public void testMinuteSurviving() throws Exception { session.commit(); - ConcurrentHashSet receivedMessages = new ConcurrentHashSet<>(); + Set receivedMessages = ConcurrentHashMap.newKeySet(); FixedRateConsumer consumer = new FixedRateConsumer(80, MESSAGES_PER_MINUTE, receivedMessages, sf, QUEUE, 0); consumer.start(); @@ -707,8 +707,8 @@ public void testMultipleConsumersOneQueue() throws Exception { FixedRateProducer producer = new FixedRateProducer(threshold * 2, MESSAGES_PER_SECOND, sf1, QUEUE, messages); - final Set consumers = new ConcurrentHashSet<>(); - final Set receivedMessages = new ConcurrentHashSet<>(); + final Set consumers = ConcurrentHashMap.newKeySet(); + final Set receivedMessages = ConcurrentHashMap.newKeySet(); consumers.add(new FixedRateConsumer(threshold, MESSAGES_PER_SECOND, receivedMessages, sf2, QUEUE, 1)); consumers.add(new FixedRateConsumer(threshold, MESSAGES_PER_SECOND, receivedMessages, sf3, QUEUE, 2)); diff --git a/tests/smoke-tests/src/test/java/org/apache/activemq/artemis/tests/smoke/failover/OpenWireSharedStoreFailoverSmokeTest.java b/tests/smoke-tests/src/test/java/org/apache/activemq/artemis/tests/smoke/failover/OpenWireSharedStoreFailoverSmokeTest.java index 79d202c8b700..842ca80e7250 100644 --- a/tests/smoke-tests/src/test/java/org/apache/activemq/artemis/tests/smoke/failover/OpenWireSharedStoreFailoverSmokeTest.java +++ b/tests/smoke-tests/src/test/java/org/apache/activemq/artemis/tests/smoke/failover/OpenWireSharedStoreFailoverSmokeTest.java @@ -24,6 +24,8 @@ import javax.jms.TextMessage; import java.io.File; import java.lang.invoke.MethodHandles; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.CyclicBarrier; import java.util.concurrent.ExecutorService; @@ -37,7 +39,6 @@ import org.apache.activemq.artemis.tests.util.CFUtil; import org.apache.activemq.artemis.util.ServerUtil; import org.apache.activemq.artemis.utils.Wait; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -127,7 +128,7 @@ public void testOpenWire() throws Exception { CyclicBarrier startFlag = new CyclicBarrier(PRODUCERS); - ConcurrentHashSet duplicateIDs = new ConcurrentHashSet<>(); + Set duplicateIDs = ConcurrentHashMap.newKeySet(); for (int producerID = 0; producerID < PRODUCERS; producerID++) { final int theProducerID = producerID; diff --git a/tests/soak-tests/src/test/java/org/apache/activemq/artemis/tests/soak/paging/ValidatePageTXTest.java b/tests/soak-tests/src/test/java/org/apache/activemq/artemis/tests/soak/paging/ValidatePageTXTest.java index 49e701bdaaa7..42d4d50af7f2 100644 --- a/tests/soak-tests/src/test/java/org/apache/activemq/artemis/tests/soak/paging/ValidatePageTXTest.java +++ b/tests/soak-tests/src/test/java/org/apache/activemq/artemis/tests/soak/paging/ValidatePageTXTest.java @@ -29,6 +29,8 @@ import javax.transaction.xa.Xid; import java.io.File; import java.lang.invoke.MethodHandles; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.CyclicBarrier; import java.util.concurrent.ExecutorService; @@ -41,7 +43,6 @@ import org.apache.activemq.artemis.cli.commands.helper.HelperCreate; import org.apache.activemq.artemis.tests.soak.SoakTestBase; import org.apache.activemq.artemis.tests.util.CFUtil; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; @@ -119,7 +120,7 @@ public void testValidatePageTX() throws Exception { AtomicInteger sequenceGenerator = new AtomicInteger(1); - ConcurrentHashSet dupList = new ConcurrentHashSet<>(); + Set dupList = ConcurrentHashMap.newKeySet(); CountDownLatch enoughSent = new CountDownLatch(SEND_GROUPS * GROUP_SIZE * 2); CountDownLatch latchDone = new CountDownLatch(SEND_GROUPS * GROUP_SIZE); @@ -258,7 +259,7 @@ private void sender(boolean useXA, AtomicInteger sequenceGenerator, AtomicBoolean running, String messageBody, - ConcurrentHashSet dupList, + Set dupList, CountDownLatch enoughSent, CountDownLatch latchDone) { try { diff --git a/tests/unit-tests/src/test/java/org/apache/activemq/artemis/tests/unit/core/journal/impl/BatchCommitTest.java b/tests/unit-tests/src/test/java/org/apache/activemq/artemis/tests/unit/core/journal/impl/BatchCommitTest.java index 9c5bafbeb819..4ec0afabe497 100644 --- a/tests/unit-tests/src/test/java/org/apache/activemq/artemis/tests/unit/core/journal/impl/BatchCommitTest.java +++ b/tests/unit-tests/src/test/java/org/apache/activemq/artemis/tests/unit/core/journal/impl/BatchCommitTest.java @@ -25,6 +25,8 @@ import java.lang.reflect.Method; import java.util.ArrayList; import java.util.List; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -47,7 +49,6 @@ import org.apache.activemq.artemis.utils.ExecutorFactory; import org.apache.activemq.artemis.utils.SimpleIDGenerator; import org.apache.activemq.artemis.utils.actors.OrderedExecutorFactory; -import org.apache.activemq.artemis.utils.collections.ConcurrentHashSet; import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; import org.slf4j.Logger; @@ -99,7 +100,7 @@ public Journal testRun(String testFolder, JournalType journalType, boolean sync) CountDownLatch latch = new CountDownLatch(RECORDS); - ConcurrentHashSet existingRecords = new ConcurrentHashSet<>(); + Set existingRecords = ConcurrentHashMap.newKeySet(); AtomicInteger errors = new AtomicInteger(0);