# UNDF: UNDF-2026-000000617 --- a/activemq-broker/src/main/java/org/apache/activemq/broker/region/Topic.java +++ b/activemq-broker/src/main/java/org/apache/activemq/broker/region/Topic.java @@ -21,6 +21,8 @@ import java.util.ArrayList; import java.util.LinkedList; import java.util.List; +import java.util.Collections; +import java.util.Set; import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.ConcurrentHashMap; @@ -83,6 +85,8 @@ public class Topic extends BaseDestination implements Task { protected final CopyOnWriteArrayList consumers = new CopyOnWriteArrayList(); + // O(1) membership guard — parallel to consumers list; updated under consumers monitor + private final Set consumerSet = + Collections.newSetFromMap(new ConcurrentHashMap()); @@ -149,7 +153,7 @@ public class Topic extends BaseDestination implements Task { boolean applyRecovery = false; synchronized (consumers) { - if (!consumers.contains(sub)){ + if (consumerSet.add(sub)){ sub.add(context, this); consumers.add(sub); applyRecovery=true; @@ -165,7 +169,7 @@ public class Topic extends BaseDestination implements Task { } else { synchronized (consumers) { - if (!consumers.contains(sub)){ + if (consumerSet.add(sub)){ sub.add(context, this); consumers.add(sub); super.addSubscription(context, sub); @@ -207,12 +211,14 @@ public class Topic extends BaseDestination implements Task { synchronized (consumers) { - removed = consumers.remove(sub); + removed = consumers.remove(sub); + if (removed) { + consumerSet.remove(sub); + } } @@ -225,6 +231,7 @@ public class Topic extends BaseDestination implements Task { synchronized (consumers) { consumers.remove(removed); + consumerSet.remove(removed); } @@ -290,7 +297,7 @@ public class Topic extends BaseDestination implements Task { synchronized (consumers) { consumers.remove(subscription); + consumerSet.remove(subscription); } } synchronized (consumers) { - if (!consumers.contains(subscription)) { + if (consumerSet.add(subscription)) { consumers.add(subscription); }