Spec-Zone.ru › OpenJDK 8

Пакет java.util.stream

Классы для поддержки операций в функциональном стиле над потоками элементов, таких как преобразования map-reduce над коллекциями.

См.: Описание

Интерфейс Описание
BaseStream<T,S extends BaseStream<T,S>>

Базовый интерфейс для потоков, которые представляют собой последовательности элементов, поддерживающие последовательные и параллельные агрегированные операции.

Collector<T,A,R>

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

DoubleStream

Последовательность примитивных элементов с двойной точностью, поддерживающая последовательные и параллельные агрегированные операции.

DoubleStream.Builder

Изменяемый билдер для DoubleStream.

IntStream

Последовательность примитивных элементов целого типа int, поддерживающая последовательные и параллельные агрегированные операции.

IntStream.Builder

Изменяемый билдер для IntStream.

LongStream

Последовательность примитивных элементов long, поддерживающая последовательные и параллельные агрегированные операции.

LongStream.Builder

Изменяемый билдер для LongStream.

Stream<T>

Последовательность элементов, поддерживающая последовательные и параллельные агрегированные операции.

Stream.Builder<T>

Изменяемый билдер для Stream.

Класс Описание
Collectors

Реализации Collector, которые реализуют различные полезные операции сокращения, такие как накопление элементов в коллекции, обобщение элементов по различным критериям и т.д.

StreamSupport

Методы низкого уровня для создания и управления потоками.

Перечисление Описание
Collector.Characteristics

Характеристики, указывающие свойства Collector, которые могут использоваться для оптимизации реализаций сокращения.

Описание пакета java.util.stream

Классы для поддержки операций в функциональном стиле над потоками элементов, таких как преобразования map-reduce над коллекциями. Например:

int sum = widgets.stream()
                      .filter(b -> b.getColor() == RED)
                      .mapToInt(b -> b.getWeight())
                      .sum();

Здесь мы используем widgets, Collection<Widget>, в качестве источника для потока, а затем выполняем фильтрацию, отображение и сокращение над потоком, чтобы получить сумму весов красных виджетов. (Суммирование является примером операции сокращения.)

Ключевой абстракцией, представленной в этом пакете, является поток. Классы Stream, IntStream, LongStream и DoubleStream являются потоками над объектами и примитивными типами int, long и double. Потоки отличаются от коллекций по нескольким аспектам:

  • Отсутствие хранения. Поток не является структурой данных, которая хранит элементы; вместо этого он передает элементы из источника, такого как структура данных, массив, функция-генератор или канал ввода-вывода, через цепочку вычислительных операций.
  • Функциональная природа. Операция над потоком производит результат, но не изменяет исходный источник. Например, фильтрация Stream, полученного из коллекции, создаёт новый Stream без отфильтрованных элементов, а не удаляет элементы из исходной коллекции.
  • Отложенное выполнение. Многие операции над потоком, такие как фильтрация, отображение или удаление дубликатов, могут быть реализованы отложенно, предоставляя возможности для оптимизации. Например, «найти первый String с тремя последовательными гласными» не требует проверки всех входных строк. Операции над потоком делятся на промежуточные (производящие Stream-поток) и терминальные (производящие значение или побочные эффекты) операции. Промежуточные операции всегда отложенные.
  • Возможность неограниченной длины. Хотя коллекции имеют конечный размер, потоки могут быть и бесконечными. Операции с ограничением, такие как limit(n) или findFirst() могут позволить вычисления над бесконечными потоками завершиться за конечное время.
  • Однократное использование. Элементы потока посещаются только один раз за время жизни потока. Как и Iterator, для повторного посещения одних и тех же элементов источника необходимо создать новый поток.
Потоки могут быть получены различными способами. Вот некоторые примеры:
  • Из Collection с помощью методов stream() и parallelStream();
  • Из массива с помощью Arrays.stream(Object[]);
  • Из статических фабричных методов классов потоков, таких как Stream.of(Object[]), IntStream.range(int, int) или Stream.iterate(Object, UnaryOperator);
  • Строки файла могут быть получены из BufferedReader.lines();
  • Потоки путей к файлам могут быть получены из методов в Files;
  • Потоки случайных чисел могут быть получены из Random.ints();
  • Многие другие методы, возвращающие потоки в JDK, включая BitSet.stream(), Pattern.splitAsStream(java.lang.CharSequence) и JarFile.stream().
Дополнительные источники потоков могут быть предоставлены сторонними библиотеками, используя эти методы.

Операции и цепочки потоков

Операции над потоками делятся на промежуточные и терминальные операции и комбинируются для формирования цепочек потоков. Цепочка потоков состоит из источника (например, Collection, массива, функции-генератора или канала ввода-вывода); за которым следуют ноль или более промежуточных операций, таких как Stream.filter или Stream.map; и терминальная операция, например, Stream.forEach или Stream.reduce.

Промежуточные операции возвращают новый поток. Они всегда отложенные; выполнение промежуточной операции, такой как filter() фактически не выполняет никакой фильтрации, а вместо этого создает новый поток, который при прохождении содержит элементы исходного потока, соответствующие заданному предикату. Перебор источника потока не начинается до выполнения терминальной операции в цепочке.

Терминальные операции, такие как Stream.forEach или IntStream.sum, могут проходить по потоку для производства результата или побочного эффекта. После выполнения терминальной операции, цепочка потоков считается использованной и больше не может быть использована; если вам нужно повторно проходить по тому же источнику данных, вы должны получить новый поток от источника данных. В большинстве случаев терминальные операции немедленные, завершая свой проход по источнику данных и обработку цепочки до возврата. Только терминальные операции iterator() и spliterator() не являются таковыми; они предоставляются как «выход» для возможности произвольного управления клиентом переходом по цепочке в случае, если существующих операций недостаточно для задачи.

Обработка потоков отложенно позволяет существенно повысить эффективность; в цепочке, такой как пример фильтрации-отображения-суммирования выше, фильтрация, отображение и суммирование могут быть объединены в один проход по данным с минимальным промежуточным состоянием. Отложенность также позволяет избегать проверки всех данных, когда это не требуется; для операций, таких как «найти первую строку, длина которой больше 1000 символов», достаточно проверить только достаточно строк, чтобы найти одну с нужными характеристиками, не проверяя все доступные строки из источника. (Это поведение становится еще более важным, когда входной поток бесконечен, а не просто большой.)

Промежуточные операции далее делятся на бессостоятельные и состоятельные операции. Бессостоятельные операции, такие как filter и map, не сохраняют состояние от ранее обработанного элемента при обработке нового элемента — каждый элемент может быть обработан независимо от операций над другими элементами. Состоятельные операции, такие как distinct и sorted, могут использовать состояние от ранее обработанных элементов при обработке новых элементов.

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

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

Параллелизм

Обработка элементов с явным for- циклом по своей природе последовательная. Потоки облегчают параллельное выполнение, переформулируя вычисления как цепочку агрегированных операций, а не как императивные операции над каждым отдельным элементом. Все операции над потоками могут выполняться как последовательно, так и параллельно. Реализации потоков в JDK создают последовательные потоки, если явное указание на параллельность не задано. Например, Collection имеет методы Collection.stream() и Collection.parallelStream(), которые производят последовательные и параллельные потоки соответственно; другие методы, возвращающие потоки, такие как IntStream.range(int, int), производят последовательные потоки, но эти потоки могут быть эффективно паралелизованы вызовом их метода BaseStream.parallel(). Для выполнения предыдущего запроса "сумма весов виджетов" параллельно, мы выполнили бы:

int sumOfWeights = widgets.parallelStream()
                               .filter(b -> b.getColor() == RED)
                               .mapToInt(b -> b.getWeight())
                               .sum();

Единственное отличие между последовательными и параллельными версиями этого примера заключается в создании начального потока, используя «parallelStream()» вместо «stream()». Когда терминальная операция инициируется, цепочка потоков выполняется последовательно или параллельно, в зависимости от ориентации потока, на котором она вызывается. Ориентацию потока можно определить с помощью метода isParallel(), и ориентацию потока можно изменить с помощью операций BaseStream.sequential() и BaseStream.parallel(). При инициировании терминальной операции цепочка потоков выполняется последовательно или параллельно в зависимости от режима потока, на котором она вызывается.

За исключением операций, явно определенных как недетерминированные, такие как findAny(), порядок выполнения потока (последовательно или параллельно) не должен изменять результат вычисления.

Большинство операций над потоками принимают параметры, которые описывают поведение, указанное пользователем, которые часто представляют собой лямбда-выражения. Для сохранения правильного поведения эти параметры поведения должны быть невмешивающимися и в большинстве случаев бессостоятельными. Такие параметры всегда являются экземплярами функционального интерфейса, такого как Function, и часто являются лямбда-выражениями или ссылками на методы.

Взаимодействие

Потоки позволяют выполнять потенциально параллельные агрегированные операции над различными источниками данных, включая даже не потокобезопасные коллекции, такие как ArrayList. Это возможно только в том случае, если мы можем предотвратить взаимодействие с источником данных во время выполнения цепочки потоков. За исключением операций «выхода», iterator() и spliterator(), выполнение начинается при вызове терминальной операции и заканчивается при её завершении. Для большинства источников данных предотвращение взаимодействия означает обеспечение того, что источник данных совсем не изменяется во время выполнения цепочки потоков. Исключением из этого являются потоки, источником которых являются потоковые коллекции, которые специально предназначены для обработки одновременных изменений. Конкурентные источники потоков — это те, чьё свойство Spliterator указывает на характеристику CONCURRENT.

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

List<String> l = new ArrayList(Arrays.asList("one", "two"));
     Stream<String> sl = l.stream();
     l.add("three");
     String s = sl.collect(joining(" "));
Сначала создаётся список, состоящий из двух строк: «one»; и «two». Затем из этого списка создаётся поток. Далее список изменяется добавлением третьей строки: «three». Наконец, элементы потока собираются и объединяются вместе. Поскольку список был изменён до начала терминальной collect операции, результатом будет строка «one two three». Все потоки, возвращаемые из коллекций JDK, и большинство других классов JDK, ведут себя таким образом; для потоков, сгенерированных другими библиотеками, см. низкоуровневое создание потоков для требований к созданию хорошо себя ведущих потоков.

Бессостоятельные поведения

Результаты потокового конвейера могут быть недетерминированными или неправильными, если поведенческие параметры операций потока являются состоятельными. Состоятельная лямбда-функция (или другой объект, реализующий соответствующий функциональный интерфейс) — это функция, результат которой зависит от любого состояния, которое может измениться во время выполнения потокового конвейера. Примером состоятельной лямбда-функции является параметр для map() в:
Set<Integer> seen = Collections.synchronizedSet(new HashSet<>());
     stream.parallel().map(e -> { if (seen.add(e)) return 0; else return e; })...
Здесь, если операция отображения выполняется параллельно, результаты для одного и того же входного значения могут отличаться от выполнения к выполнению из-за различий в планировании потоков, тогда как со статической лямбда-выражением результаты всегда будут одинаковыми.

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

Побочные эффекты

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

Если поведенческие параметры имеют побочные эффекты, за исключением случаев явного указания, нет гарантий относительно видимости этих побочных эффектов для других потоков, а также нет гарантий, что разные операции над «одним» элементом в одном и том же потоковом конвейере выполняются в одном и том же потоке. Кроме того, порядок этих эффектов может быть неожиданным. Даже когда конвейер ограничен для получения результата, согласованного с порядком встречи источника потока (например, IntStream.range(0,5).parallel().map(x -> x*2).toArray() должен производить [0, 2, 4, 6, 8]), не даются никакие гарантии относительно порядка, в котором функция отображения применяется к отдельным элементам, или в каком потоке выполняется любой поведенческий параметр для данного элемента.

Многие вычисления, в которых можно было бы использовать побочные эффекты, могут быть более безопасными и эффективными без них, например, используя сведение вместо изменяемых накопителей. Однако побочные эффекты, такие как использование println() для отладки, обычно безвредны. Небольшое количество операций потока, таких как forEach() и peek(), могут работать только через побочные эффекты; их следует использовать с осторожностью.

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

ArrayList<String> results = new ArrayList<>();
     stream.filter(s -> pattern.matcher(s).matches())
           .forEach(s -> results.add(s));  // Unnecessary use of side-effects!
Этот код необоснованно использует побочные эффекты. Если он будет выполнен параллельно, небезопасность для потоков ArrayList приведёт к неправильным результатам, а добавление необходимой синхронизации приведёт к конфликтам, подорвав преимущества параллелизма. Более того, использование побочных эффектов здесь совершенно излишне; forEach() можно просто заменить операцией сведения, которая является более безопасной, более эффективной и более подходящей для параллелизации:
List<String>results =
         stream.filter(s -> pattern.matcher(s).matches())
               .collect(Collectors.toList());  // No side-effects!

Порядок

Потоки могут или не могут иметь определенный порядок встречи. Наличие или отсутствие порядка встречи потока зависит от источника и промежуточных операций. Некоторые источники потоков (например, List или массивы) изначально упорядочены, тогда как другие (например, HashSet) нет. Некоторые промежуточные операции, такие как sorted(), могут накладывать порядок встречи на поток, который в противном случае неупорядочен, а другие могут сделать упорядоченный поток неупорядоченным, например, BaseStream.unordered(). Кроме того, некоторые терминальные операции могут игнорировать порядок встречи, например, forEach().

Если поток упорядочен, большинство операций ограничены выполнением над элементами в порядке их встречи; если источником потока является List содержащий [1, 2, 3], то результат выполнения map(x -> x*2) должен быть [2, 4, 6]. Однако, если у источника нет определённого порядка встречи, любая перестановка значений [2, 4, 6] будет допустимым результатом.

Для последовательных потоков наличие или отсутствие порядка встречи не влияет на производительность, а только на детерминизм. Если поток упорядочен, повторное выполнение идентичных потоковых конвейеров на идентичном источнике даст идентичный результат; если он неупорядочен, повторное выполнение может привести к различным результатам.

Для параллельных потоков ослабление ограничения порядка иногда может позволить более эффективное выполнение. Некоторые агрегирующие операции, такие как фильтрация дубликатов (distinct()) или групповые сводки (Collectors.groupingBy()) могут быть реализованы более эффективно, если порядок элементов не важен. Аналогично, операции, которые тесно связаны с порядком встречи, такие как limit(), могут потребовать буферизации для обеспечения надлежащего порядка, подорвав преимущества параллелизма. В тех случаях, когда у потока есть порядок встречи, но пользователю он особенно не нужен, явное отключение порядка потока с помощью unordered() может улучшить производительность параллелизма для некоторых состоятельных или терминальных операций. Тем не менее, большинство потоковых конвейеров, таких как пример «суммы веса блоков», по-прежнему эффективно параллелизуются даже с ограничениями на порядок.

Операции сведения

Операция сведения (также называемая сведением) принимает последовательность входных элементов и объединяет их в один итоговый результат путём многократного применения операции объединения, например, нахождения суммы или максимума набора чисел или накопления элементов в список. Классы потоков имеют несколько форм общих операций сведения, называемых reduce() и collect(), а также несколько специализированных форм сведения, таких как sum(), max() или count().

Конечно, такие операции можно легко реализовать в виде простых последовательных циклов, как в:

int sum = 0;
    for (int x : numbers) {
       sum += x;
    }
Однако существуют веские причины предпочесть операцию сведения итеративной накапливающей операции, такой как показанная выше. Операция сведения не только «более абстрактна» — она работает со всем потоком, а не с отдельными элементами, — но и правильно построенная операция сведения по своей природе параллелизуема, если функции обработки элементов являются ассоциативными и бессостоятельными. Например, для нахождения суммы потока чисел можно записать:
int sum = numbers.stream().reduce(0, (x,y) -> x+y);
или:
int sum = numbers.stream().reduce(0, Integer::sum);

Эти операции сведения могут выполняться безопасно в параллельном режиме практически без модификаций:

int sum = numbers.parallelStream().reduce(0, Integer::sum);

Сведение хорошо параллелизуется, потому что реализация может работать с подмножествами данных в параллельном режиме, а затем комбинировать промежуточные результаты, чтобы получить окончательный правильный ответ. (Даже если язык имел конструкцию «параллельный цикл foreach», подходу итеративной накапливающей операции всё равно потребовалось бы предоставление разработчиком потокобезопасных обновлений для общей переменной накопления sum, и необходимая синхронизация, вероятно, устранила бы любые преимущества параллелизма.) Использование reduce() вместо этого устраняет всю нагрузку по параллелизации операции сведения, и библиотека может предоставить эффективную параллельную реализацию без дополнительной синхронизации.

Примеры "виджеты", показанные ранее, демонстрируют, как сведение сочетается с другими операциями для замены циклов for операциями пакетной обработки. Если widgets представляет собой коллекцию Widget объектов, у которых есть метод getWeight, мы можем найти самый тяжёлый виджет с помощью:

OptionalInt heaviest = widgets.parallelStream()
                                   .mapToInt(Widget::getWeight)
                                   .max();

В своей наиболее общей форме операция reduce над элементами типа <T>, дающая результат типа <U>, требует трёх параметров:

<U> U reduce(U identity,
              BiFunction<U, ? super T, U> accumulator,
              BinaryOperator<U> combiner);
Здесь идентичный элемент является одновременно начальным значением семени для сведения и значением по умолчанию для результата, если входных элементов нет. Функция накопления принимает частичный результат и следующий элемент и создаёт новый частичный результат. Функция комбинирования объединяет два частичных результата, чтобы создать новый частичный результат. (Функция комбинирования необходима при параллельном сведении, когда вход разделяется, для каждой части вычисляется частичное накопление, а затем частичные результаты объединяются для получения окончательного результата.)

Более формально, значение identity должно быть идентичным для функции комбинирования. Это означает, что для всех u, combiner.apply(identity, u) равно u. Кроме того, функция combiner должна быть ассоциативной и должна быть совместима с функцией accumulator: для всех u и t, combiner.apply(u, accumulator.apply(identity, t)) должно быть equals() к accumulator.apply(u, t).

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

int sumOfWeights = widgets.stream()
                               .reduce(0,
                                       (sum, b) -> sum + b.getWeight())
                                       Integer::sum);
хотя явная форма map-reduce более удобочитаема и поэтому обычно предпочтительнее. Обобщенная форма предоставляется для случаев, когда существенную работу можно оптимизировать, объединив отображение и сокращение в одну функцию.

Изменяемое сокращение

Операция изменяемого сокращения накапливает элементы входных данных в контейнер изменяемого результата, такой как Collection или StringBuilder, при обработке элементов в потоке.

Если мы захотели бы взять поток строк и конкатенировать их в одну длинную строку, мы могли бы добиться этого с помощью обычного сокращения:

String concatenated = strings.reduce("", String::concat)

Мы бы получили желаемый результат, и он даже работал бы параллельно. Однако нас, возможно, не порадовала бы производительность! Такая реализация выполнит значительное количество копирования строк, и время выполнения будет O(n^2) по количеству символов. Более производительный подход заключался бы в накоплении результатов в StringBuilder, который является изменяемым контейнером для накопления строк. Мы можем использовать ту же технику для паралелизации изменяемого сокращения, что и для обычного сокращения.

Операция изменяемого сокращения называется collect(), поскольку она собирает вместе желаемые результаты в контейнер результата, такой как Collection. Операция collect требует трех функций: функцию поставщика для построения новых экземпляров контейнера результата, функцию накопителя для включения элемента входных данных в контейнер результата и функцию объединения для слияния содержимого одного контейнера результата в другой. Форма очень похожа на общую форму обычного сокращения:

<R> R collect(Supplier<R> supplier,
               BiConsumer<R, ? super T> accumulator,
               BiConsumer<R, R> combiner);

Как и в случае с reduce(), преимущество выражения collect таким абстрактным способом заключается в том, что оно непосредственно подходит для паралелизации: мы можем накапливать частичные результаты параллельно, а затем объединять их, при условии, что функции накопления и объединения удовлетворяют соответствующим требованиям. Например, чтобы собрать строковые представления элементов в потоке в ArrayList, мы могли бы написать очевидную последовательную форму for-each:

ArrayList<String> strings = new ArrayList<>();
     for (T element : stream) {
         strings.add(element.toString());
     }
Или мы могли бы использовать паралелизуемую форму collect:
ArrayList<String> strings = stream.collect(() -> new ArrayList<>(),
                                                (c, e) -> c.add(e.toString()),
                                                (c1, c2) -> c1.addAll(c2));
или, вынося операцию отображения из функции накопителя, мы могли бы выразить её более лаконично как:
List<String> strings = stream.map(Object::toString)
                                  .collect(ArrayList::new, ArrayList::add, ArrayList::addAll);
Здесь наш поставщик — просто ArrayList constructor, накопитель добавляет строковое представление элемента в ArrayList, а объединитель просто использует addAll для копирования строк из одного контейнера в другой.

Три аспекта collect — поставщик, накопитель и объединитель — тесно связаны. Мы можем использовать абстракцию Collector для захвата всех трёх аспектов. Приведенный выше пример сбора строк в List можно переписать с использованием стандартного Collector как:

List<String> strings = stream.map(Object::toString)
                                  .collect(Collectors.toList());

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

Collector<Employee, ?, Integer> summingSalaries
         = Collectors.summingInt(Employee::getSalary);
(? для второго параметра типа просто указывает, что нас не интересует промежуточное представление, используемое этим коллектором.) Если мы хотели бы создать коллектор для таблицы сумм зарплат по отделам, мы могли бы повторно использовать summingSalaries с помощью groupingBy:
Map<Department, Integer> salariesByDept
         = employees.stream().collect(Collectors.groupingBy(Employee::getDepartment,
                                                            summingSalaries));

Как и в случае с обычной операцией сокращения, операции collect() могут быть паралелизованы только при выполнении соответствующих условий. Для любого частично накопленного результата объединение его с пустым контейнером результата должно давать эквивалентный результат. То есть, для частично накопленного результата p, являющегося результатом любой серии вызовов накопителя и объединителя, p должен быть эквивалентен combiner.apply(p, supplier.get()).

Кроме того, независимо от того, как вычисление разбивается, оно должно давать эквивалентный результат. Для любых элементов входных данных t1 и t2, результаты r1 и r2 в вычислении ниже должны быть эквивалентны:

A a1 = supplier.get();
     accumulator.accept(a1, t1);
     accumulator.accept(a1, t2);
     R r1 = finisher.apply(a1);  // result without splitting

     A a2 = supplier.get();
     accumulator.accept(a2, t1);
     A a3 = supplier.get();
     accumulator.accept(a3, t2);
     R r2 = finisher.apply(combiner.apply(a2, a3));  // result with splitting

Здесь эквивалентность обычно означает в соответствии с Object.equals(Object), но в некоторых случаях эквивалентность может быть ослаблена для учета различий в порядке.

Сокращение, конкуреность и порядок

При некоторых сложных операциях сокращения, например, collect() , которые производят Map, например:
Map<Buyer, List<Transaction>> salesByBuyer
         = txns.parallelStream()
               .collect(Collectors.groupingBy(Transaction::getBuyer));
, параллельное выполнение операции может быть на самом деле неэффективным. Это происходит потому, что этап объединения (слияние одного Map в другой по ключу) может быть дорогим для некоторых Map реализаций.

Однако, предположим, что контейнер результата, используемый в этом сокращении, был одновременно изменяемым коллекцией — такой как ConcurrentHashMap. В этом случае параллельные вызовы накопителя фактически могли бы одновременно помещать свои результаты в один и тот же общий контейнер результата, устраняя необходимость для объединителя в слиянии различных контейнеров результатов. Это потенциально улучшает производительность параллельного выполнения. Мы называем это конкурентным сокращением.

Collector, который поддерживает конкурентное сокращение, помечен характеристикой Collector.Characteristics.CONCURRENT. Однако у конкурентного коллектора также есть недостаток. Если несколько потоков одновременно помещают результаты в общий контейнер, порядок, в котором результаты размещаются, не определён. Следовательно, конкурентное сокращение возможно только в том случае, если порядок не важен для обрабатываемого потока. Реализация Stream.collect(Collector) выполнит конкурентное сокращение только в том случае, если

  • Поток является параллельным;
  • Коллектор имеет характеристику Collector.Characteristics.CONCURRENT, и;
  • Либо поток не упорядочен, либо коллектор имеет характеристику Collector.Characteristics.UNORDERED.
Вы можете убедиться, что поток не упорядочен, используя метод BaseStream.unordered(). Например:
Map<Buyer, List<Transaction>> salesByBuyer
         = txns.parallelStream()
               .unordered()
               .collect(groupingByConcurrent(Transaction::getBuyer));
(где Collectors.groupingByConcurrent(java.util.function.Function<? super T, ? extends K>) — это конкурентный эквивалент groupingBy).

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

Ассоциативность

Оператор или функция op является ассоциативной, если выполняется следующее условие:
(a op b) op c == a op (b op c)
Важность этого для параллельной оценки можно увидеть, если мы расширим это до четырех терминов:
a op b op c op d == (a op b) op (c op d)
Итак, мы можем оценить (a op b) параллельно с (c op d), а затем вызвать op на результатах.

Примеры ассоциативных операций включают числовое сложение, мин и макс, и конкатенацию строк.

Построение потоков низкого уровня

До сих пор все примеры потоков использовали методы, такие как Collection.stream() или Arrays.stream(Object[]) для получения потока. Как реализуются эти методы, порождающие потоки?

В классе StreamSupport есть ряд методов низкого уровня для создания потока, все они используют какую-то форму Spliterator. Spliterator — это параллельный аналог Iterator; он описывает (возможно, бесконечный) набор элементов с поддержкой последовательного продвижения, массового обхода и разделения части входных данных на другой spliterator, который может обрабатываться параллельно. На самом низком уровне все потоки управляются spliterator.

Существует ряд вариантов реализации spliterator, почти все из которых представляют собой компромиссы между простотой реализации и производительностью потоков, использующих этот spliterator. Наиболее простой, но наименее производительный способ создания spliterator — создание его из итератора с помощью Spliterators.spliteratorUnknownSize(java.util.Iterator, int). Хотя такой spliterator будет работать, он, вероятно, предложит плохую параллельную производительность, поскольку мы потеряли информацию о размере (какой размер имеет основной набор данных), а также ограничены простым алгоритмом разделения.

Spliterator более высокого качества обеспечит сбалансированные и известные разделения, точную информацию о размере и ряд других characteristics spliterator или данных, которые могут быть использованы реализациями для оптимизации выполнения.

Для изменяемых источников данных spliterator имеют дополнительную проблему; момент привязки к данным, так как данные могут измениться между моментом создания spliterator и моментом выполнения потоковой обработки. В идеале spliterator для потока должен сообщать о характеристике IMMUTABLE или CONCURRENT; если нет, то он должен быть поздневязанным. Если источник не может напрямую предоставить рекомендуемый spliterator, он может косвенно предоставить spliterator, используя Supplier, и создать поток с помощью версий stream(), принимающих Supplier. Spliterator извлекается из поставщика только после начала терминальной операции потока.

Эти требования значительно сужают область потенциальных конфликтов между изменениями источника потока и выполнением потоковых конвейеров. Потоки, основанные на spliterator с требуемыми характеристиками или тех, которые используют формы фабрики на основе Supplier, устойчивы к изменениям источника данных до начала терминальной операции (при условии, что параметры поведения потоковых операций удовлетворяют необходимым критериям для отсутствия конфликтов и бессостоятельности). Подробности см. в разделе Независимость.

С момента:
1.8

© 1993, 2020, Oracle and/or its affiliates. All rights reserved.
Documentation extracted from Debian's OpenJDK Development Kit package.
Licensed under the GNU General Public License, version 2, with the Classpath Exception.
Various third party code in OpenJDK is licensed under different licenses (see Debian package).
Java and OpenJDK are trademarks or registered trademarks of Oracle and/or its affiliates.

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API