寻求有序限容的线程安全集合:每分钟从10万数据筛选Top20
一、JDK自带实现(无需额外依赖)
1. 基于PriorityQueue+显式锁
用PriorityQueue维护最小堆,只保留Top N元素,结合ReentrantLock保证线程安全,适合高并发场景:
import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.PriorityQueue; import java.util.concurrent.locks.ReentrantLock; public class BoundedSortedTopSet<T extends Comparable<T>> { private final int maxSize; private final PriorityQueue<T> queue; private final ReentrantLock lock = new ReentrantLock(); public BoundedSortedTopSet(int maxSize) { this.maxSize = maxSize; // 最小堆,堆顶是当前集合中最小的元素 this.queue = new PriorityQueue<>(maxSize); } public void add(T element) { lock.lock(); try { if (queue.size() < maxSize) { queue.offer(element); } else { // 仅当新元素比堆顶大时,替换堆顶 if (element.compareTo(queue.peek()) > 0) { queue.poll(); queue.offer(element); } } } finally { lock.unlock(); } } // 获取从大到小排序的Top N元素 public List<T> getTopElements() { lock.lock(); try { List<T> list = new ArrayList<>(queue); list.sort(Collections.reverseOrder()); return Collections.unmodifiableList(list); } finally { lock.unlock(); } } }
这个实现的插入操作时间复杂度为O(logN),锁粒度小,能高效处理每分钟10万级别的数据插入。
2. 基于ConcurrentSkipListSet+容量控制
ConcurrentSkipListSet是线程安全的有序集合,手动控制容量,插入后检查大小,超过则删除最小元素:
import java.util.ArrayList; import java.util.List; import java.util.concurrent.ConcurrentSkipListSet; public class BoundedConcurrentSortedSet<T extends Comparable<T>> { private final int maxSize; private final ConcurrentSkipListSet<T> set; public BoundedConcurrentSortedSet(int maxSize) { this.maxSize = maxSize; this.set = new ConcurrentSkipListSet<>(); } public boolean add(T element) { boolean added = set.add(element); if (added && set.size() > maxSize) { // 删除当前集合中最小的元素 set.pollFirst(); } return added; } public List<T> getTopElements() { // 返回从大到小排序的元素集合 return new ArrayList<>(set.descendingSet()); } }
注意:add和size()检查并非原子操作,极端情况下可能短暂出现容量超标的情况,但后续的pollFirst会自动修正。如果需要严格保证容量不超限,可以用显式锁包裹add与容量检查的逻辑。
二、第三方库实现
1. Guava的MinMaxPriorityQueue
Guava提供的MinMaxPriorityQueue支持快速获取最小/最大元素,可直接设置固定容量,超过时自动删除最小元素:
import com.google.common.collect.MinMaxPriorityQueue; import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Collections; public class GuavaBoundedTopSet<T extends Comparable<T>> { private final MinMaxPriorityQueue<T> queue; public GuavaBoundedTopSet(int maxSize) { this.queue = MinMaxPriorityQueue.maximumSize(maxSize).create(); } public void add(T element) { queue.add(element); // 超过容量自动删除最小元素 } public List<T> getTopElements() { List<T> list = new ArrayList<>(queue); list.sort(Collections.reverseOrder()); return Collections.unmodifiableList(list); } }
注意:MinMaxPriorityQueue本身非线程安全,需用Collections.synchronizedCollection包装:
MinMaxPriorityQueue<T> queue = MinMaxPriorityQueue.maximumSize(maxSize).create(); Collection<T> synchronizedQueue = Collections.synchronizedCollection(queue);
2. Apache Commons Collections的BoundedSortedSet
Apache Commons Collections 4.x的BoundedSortedSet是固定容量的有序集合,可配置超过容量时删除最小元素或拒绝添加:
import org.apache.commons.collections4.set.BoundedSortedSet; import java.util.ArrayList; import java.util.List; import java.util.TreeSet; import java.util.Collections; public class CommonsBoundedTopSet<T extends Comparable<T>> { private final BoundedSortedSet<T> boundedSet; public CommonsBoundedTopSet(int maxSize) { // 第三个参数为false时,超容删除最小元素;为true时拒绝添加新元素 this.boundedSet = BoundedSortedSet.boundedSortedSet(new TreeSet<>(), maxSize, false); } public boolean add(T element) { return boundedSet.add(element); } public List<T> getTopElements() { return new ArrayList<>(boundedSet.descendingSet()); } }
注意:BoundedSortedSet本身非线程安全,需用Collections.synchronizedSortedSet包装以保证并发安全。
三、CopyOnWriteArrayList异常原因
CopyOnWriteArrayList的所有修改操作都会复制底层数组,迭代器基于修改前的数组快照。在并发场景下执行以下逻辑时:
list.add(element); if (list.size() > 20) { list.remove(20); }
可能在add后、size()检查前,其他线程已执行add操作,导致size()返回值与实际集合大小不一致。当执行remove(20)时,底层数组已被修改,索引20可能不存在,从而抛出ConcurrentModificationException。
总结
- 无第三方库依赖时,优先选择
PriorityQueue+ReentrantLock的实现,性能最优,适配高并发场景; - 允许引入第三方库时,Guava或Apache Commons的封装类可减少自定义代码,但需注意线程安全包装;
- 避免使用
CopyOnWriteArrayList处理此类需求,其快照特性无法满足实时一致性的容量控制要求。
内容的提问来源于stack exchange,提问作者Oleksii Migunov

