You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

寻求有序限容的线程安全集合:每分钟从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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.29 14:43:19