Напишите свой Collector, оставляющий n наибольших элементов
Реализуйте topN — переиспользуемый Collector, сворачивающий stream до n наибольших элементов по заданному Comparator, от большего к меньшему. Соберите его через Collector.of — не сводите задачу к sorted().limit(n).
Ограничения:
- результат должен совпадать на параллельном stream, поэтому нужен корректный combiner
- держать одновременно не более
nэлементов; не сортировать и не буферизовать весь вход nбольше числа элементов возвращает их все;n <= 0возвращает пустой список
static <T> Collector<T, ?, List<T>> topN(int n, Comparator<? super T> cmp) {
// ваш код здесь
}
Допишите реализацию.
Строит его Collector.of(supplier, accumulator, combiner, finisher). Supplier создаёт свежий изменяемый контейнер на воркера, аккумулятор вкладывает один элемент, комбайнер сливает два контейнера — это и делает коллектор безопасным в параллели, — а finisher превращает контейнер в результат. Для topN контейнер — min-heap: кладём элемент, выбрасываем наименьший, как только их больше n, и сортируем уцелевших в finisher. Память остаётся O(n).
- ✗Писать комбайнер, теряющий частичный результат одной из сторон, из-за чего параллельный ответ расходится с последовательным
- ✗Буферизовать все элементы и сортировать в finisher, из-за чего память становится O(входа), а не O(n)
- ✗Возвращать сам изменяемый контейнер вместо его преобразования в finisher
- →Что обещает характеристика
UNORDEREDи когда её безопасно объявлять? - →Когда применима
IDENTITY_FINISHи что она экономит конвейеру?
Решение
static <T> Collector<T, ?, List<T>> topN(int n, Comparator<? super T> cmp) {
return Collector.of(
() -> new PriorityQueue<T>(cmp), // min-heap: наименьший из отобранных — сверху
(heap, item) -> { // аккумулятор: один элемент
heap.offer(item);
if (heap.size() > n) heap.poll(); // выбрасываем наименьшего
},
(a, b) -> { // комбайнер: слияние двух частичных куч
for (T item : b) {
a.offer(item);
if (a.size() > n) a.poll();
}
return a;
},
heap -> { // finisher: куча → список
List<T> out = new ArrayList<>(heap);
out.sort(cmp.reversed());
return out;
},
Collector.Characteristics.UNORDERED);
}
// использование
List<Employee> top3 = staff.stream()
.collect(topN(3, Comparator.comparingInt(Employee::salary)));
Почему так. Контейнер — PriorityQueue с исходным cmp, то есть min-heap: наверху лежит наименьший из уже отобранных. Аккумулятор кладёт элемент и, если размер превысил n, снимает вершину. Так в памяти живёт не более n + 1 элементов — вход не сортируется целиком.
Комбайнер — не украшение. На параллельном stream каждый воркер накапливает свою кучу, и слить их обязан именно комбайнер: он вливает элементы одной кучи в другую, соблюдая тот же порог n. Комбайнер, который вернул бы только a, молча терял бы половину данных — и только под параллелью.
Края. При n <= 0 порог heap.size() > n срабатывает сразу же, поэтому куча остаётся пустой и результат — пустой список. При n больше числа элементов не выбрасывается ничего, и возвращается весь вход, отсортированный по убыванию. UNORDERED заявляет, что порядок появления элементов на результат не влияет, и разрешает рантайму не поддерживать его.