Разберите ошибку в объявлении Collector: почему параллельный сбор может повредить результат, несмотря на на...

Разберите ошибку в объявлении Collector: почему параллельный сбор может повредить результат, несмотря на наличие combiner?

import java.util.*;
import java.util.stream.*;

Collector<Integer, List<Integer>, List<Integer>> bad = Collector.of(
    ArrayList::new,
    List::add,
    (left, right) -> { left.addAll(right); return left; },
    Collector.Characteristics.CONCURRENT
);

List<Integer> result = IntStream.range(0, 100_000)
    .parallel().unordered().boxed().collect(bad);
Проходите собеседования с ИИ помощником Hintsage

Краткий ответ

Ошибка в объявлении характеристики CONCURRENT: она разрешает нескольким потокам одновременно вызывать accumulator для одного и того же контейнера. ArrayList не поддерживает безопасную конкурентную модификацию, поэтому возможны потеря элементов, повреждение внутреннего состояния или исключение.

combiner не спасает ситуацию: при конкурентном накоплении несколько частичных контейнеров могут вообще не создаваться. Исправление — убрать CONCURRENT либо использовать действительно потокобезопасный контейнер и соответствующую ему стратегию накопления.

Исторический контекст

Collector отделяет описание алгоритма сбора от его выполнения. Это позволяет Stream API выбирать последовательную или параллельную стратегию, создавая частичные результаты и объединяя их через combiner.

Характеристики коллектора сообщают реализации дополнительные свойства операции. В частности, CONCURRENT предназначена для случаев, когда один контейнер результата допускает одновременное накопление из нескольких потоков; без такого свойства реализация обычно использует независимые контейнеры и затем объединяет их.

Постановка проблемы

В примере ArrayList используется как изменяемый контейнер, а List::add — как аккумулятор. После объявления CONCURRENT параллельный неупорядоченный стрим может направить вызовы add из разных потоков в один экземпляр ArrayList.

Метод ArrayList.add не синхронизирован и не использует атомарное обновление размера и массива. Поэтому итоговый размер может быть меньше ожидаемого, элементы могут потеряться, а внутреннее состояние списка — стать некорректным.

Подробное решение

Для параллельного сбора возможны два принципиально разных сценария:

  1. Без CONCURRENT потоки накапливают данные в отдельных контейнерах. Затем Stream API вызывает combiner и объединяет эти контейнеры. Каждый ArrayList в таком сценарии обычно изменяется только одним потоком.
  2. С CONCURRENT реализация получает право использовать общий контейнер. Для неупорядоченного стрима это позволяет нескольким потокам одновременно вызывать аккумулятор для него.

В примере .unordered() снимает требование сохранять порядок элементов. Вместе с CONCURRENT это создаёт условия для общего накопления в одном ArrayList. Наличие корректного combiner не делает accumulator потокобезопасным: combiner отвечает только за объединение контейнеров, когда такое объединение выполняется.

Простейшее исправление — не объявлять ошибочную характеристику:

Collector<Integer, List<Integer>, List<Integer>> safe = Collector.of( ArrayList::new, List::add, (left, right) -> { left.addAll(right); return left; } ); List<Integer> result = IntStream.range(0, 100_000) .parallel().unordered().boxed().collect(safe);

Теперь параллельная реализация должна соблюдать обычную модель частичных результатов и их объединения. Это безопаснее, но не гарантирует ускорение: создание контейнеров, синхронизация внутренних механизмов параллельного стрима и объединение списков могут стоить дороже последовательного сбора.

Если нужен именно общий конкурентный контейнер, аккумулятор должен быть безопасен при одновременных вызовах. Например, можно использовать ConcurrentLinkedQueue, однако это меняет тип и свойства результата. CopyOnWriteArrayList обычно не подходит для массового добавления: каждое изменение может копировать внутренний массив, поэтому он предназначен главным образом для сценариев с редкими записями и частым чтением.

Важно различать CONCURRENT и потокобезопасность всего результата. Характеристика является обещанием коллектора Stream API, а не автоматической синхронизацией контейнера. Объявлять её можно только после проверки аккумулятора, требований к порядку и поведения финализатора.

Ситуация из практики

Сервис собирал большой набор идентификаторов в список и пытался ускорить обработку, добавив CONCURRENT к коллектору с ArrayList. На тестовых данных ошибка проявлялась нестабильно: иногда размер результата совпадал с ожидаемым, а иногда часть идентификаторов исчезала.

Рассматривались три варианта. Удаление CONCURRENT сохраняло тип результата и корректность, но оставляло стоимость объединения частичных списков. Замена контейнера на CopyOnWriteArrayList обеспечивала потокобезопасность, но делала массовую запись слишком дорогой. Использование конкурентной очереди подходило для параллельного накопления, но требовало изменить контракт результата и отдельно решить вопрос порядка.

Выбрали обычный коллектор без CONCURRENT, а параллельность оставили только после измерений. В результате корректность обеспечивалась моделью отдельных контейнеров, а для исходного объёма данных последовательный сбор оказался не медленнее из-за отсутствия затрат на распараллеливание и объединение.

Что кандидаты часто упускают

  1. Достаточно ли сделать ArrayList синхронизированным, чтобы всегда корректно объявить CONCURRENT?

Нет. Синхронизированный список защищает отдельные операции, но нужно проверить весь контракт коллектора, включая финальный результат, требования к порядку и совместимость с характеристиками. Кроме того, блокировка на каждой операции может устранить преимущество параллельного накопления.

Если порядок не важен и контейнер действительно допускает конкурентные вызовы аккумулятора, такое решение может быть корректным. Но это должно быть осознанным компромиссом, подтверждённым измерениями.

  1. Что изменится, если убрать .unordered() и оставить CONCURRENT?

У стрима останется требование сохранять encounter order, а коллектор не объявляет UNORDERED. Поэтому реализация не обязана использовать общий контейнер для одновременного накопления: она может выбрать раздельные контейнеры и последующее объединение, чтобы сохранить порядок.

Это не превращает ошибочный коллектор в надёжный во всех сценариях. Нельзя полагаться на конкретный выбор реализации, потому что контракт CONCURRENT всё равно сообщает, что аккумулятор допускает конкурентное использование.

  1. Почему CONCURRENT иногда не ускоряет сбор, даже если контейнер потокобезопасен?

Общий контейнер устраняет часть операций объединения, но добавляет конкуренцию за его внутренние структуры и возможные блокировки или атомарные операции. При большом числе потоков они могут стать узким местом.

Кроме того, параллельный стрим имеет накладные расходы на разбиение источника, планирование задач и координацию. Если элементы обрабатываются дёшево или данных мало, последовательный сбор часто оказывается быстрее; выбор нужно делать по измерениям на реальной нагрузке.