В параллельном стриме операция limit может запросить у источника больше элементов, чем вернёт. Каким механизмом объясняется такое поведение?
limit ограничивает результат, но не обязано ограничивать число элементов, уже затронутых параллельной обработкой. Разные задачи работают одновременно, поэтому к моменту обнаружения первых нужных элементов другие задачи могут уже извлечь и обработать дополнительные элементы.
Stream API, появившийся в Java 8, предназначен для декларативной обработки данных с возможностью автоматического распараллеливания. Для этого источник разбивается на части, а независимые задачи обрабатываются параллельно.
Операция limit относится к короткозамыкающим операциям: она может завершить вычисление, когда набрано нужное количество элементов. В последовательном потоке остановка обычно происходит сразу после получения требуемого префикса, но параллельное выполнение требует координации уже запущенных задач.
Результат limit(n) содержит не более n элементов, однако это не означает, что источник посетит ровно n элементов. Дополнительные элементы могут быть обработаны, отброшены или временно сохранены во внутренних буферах.
Это особенно важно для дорогих операций, побочных эффектов и источников с внешними ресурсами. Ошибочно считать, что limit гарантированно уменьшает стоимость параллельного вычисления до стоимости обработки только первых n элементов.
Параллельный стрим делит источник на части и передаёт их отдельным задачам. Каждая задача может начать извлекать элементы, пока другая задача ещё не определила, какие элементы входят в итоговый префикс.
Для упорядоченного стрима limit должен сохранить первые элементы согласно encounter order. Поэтому результаты разных задач приходится сопоставлять, иногда буферизовать и отбрасывать элементы, которые оказались после нужного префикса. После достижения лимита система подаёт сигнал отмены другим задачам, но уже выполняющиеся задачи не обязаны мгновенно остановиться.
Для неупорядоченного стрима можно вернуть любые n элементов. Это обычно упрощает координацию и позволяет быстрее остановить обработку, но дополнительные элементы всё равно могут быть затронуты до распространения сигнала отмены.
Минимальный пример показывает различие между размером результата и числом посещённых элементов:
Размер result равен 10, но значение visited может быть больше 10. Точное количество посещённых элементов не является контрактом, поэтому peek не следует использовать для критически важных побочных эффектов или точного подсчёта работы.
Параллельный limit особенно осторожно следует применять к бесконечным упорядоченным источникам: поиск первых элементов может требовать значительной координации и буферизации. Если порядок не нужен, unordered() иногда заметно улучшает масштабирование, но меняет допустимый набор элементов результата.
Сервис выбирает 100 записей из большого удалённого источника и применяет к каждой записи дорогую проверку. Разработчик использует параллельный стрим с limit(100) и ожидает, что проверка выполнится ровно для 100 записей.
Вариант с обычным упорядоченным параллельным стримом сохраняет порядок, но может проверять дополнительные записи из-за уже работающих задач и координации префикса. Вариант с unordered() уменьшает требования к координации, однако возвращает любые подходящие записи, а не первые по исходному порядку.
Если бизнес-требование допускает произвольные 100 записей, выбирается неупорядоченная обработка и измеряется фактическая стоимость. Если нужны именно первые 100, порядок сохраняется, но ограничение переносится ближе к источнику данных, например в запрос или API поставщика. Это лучше, чем полагаться на то, что limit остановит всю дорогую работу ровно на сотом элементе.
limit обработку ровно указанного количества элементов?Нет. Он гарантирует ограничение результата: будет возвращено не более заданного числа элементов. В параллельном стриме задачи могут обработать дополнительные элементы до того, как информация о достижении лимита распространится по вычислению.
unordered()?Стрим получает право вернуть любые элементы, а не обязательный префикс encounter order. Это может сократить ожидание между задачами и уменьшить буферизацию, поэтому limit часто масштабируется лучше. Но unordered() не обещает, что будет обработано ровно нужное количество элементов, и не сохраняет исходный порядок результата.
peek для освобождения ресурсов после достижения лимита?Нет, это ненадёжно. peek предназначен главным образом для наблюдения и отладки, а его вызовы в параллельном вычислении могут происходить для дополнительных элементов; порядок и точное количество вызовов не следует использовать как протокол управления ресурсами.
Ресурсы нужно связывать с корректным временем жизни источника и закрывать через подходящий механизм, например try-with-resources, если источник поддерживает AutoCloseable. Побочные эффекты внутри функций стрима также усложняют повторное использование, тестирование и параллельную безопасность.