помещение
Эта статья знакомит со скрытыми ловушками ParallelStream в повседневной разработке.Эта проблема на самом деле очень близка нам, особенно если мы любим использоватьJDK1.8+Партнер по потоковому программированию должен быть глубоко тронут. Так называемая «ловушка» в заголовке на самом деле не является ловушкой самого ParallelStream, но обычно является ловушкой, которую разработчики ошибочно используют ParallelStream, чтобы похоронить себя.
Преднамеренный пример
Ниже приведен намеренный пример, на самом деле подобного бизнес-кода быть не должно:
public class ParallelStreamMain {
public static void main(String[] args) throws Exception {
List<List<Integer>> array = new ArrayList<>();
List<Integer> item1 = new ArrayList<>();
List<Integer> item2 = new ArrayList<>();
List<Integer> target = new ArrayList<>(100);
array.add(item1);
array.add(item2);
array.parallelStream().forEach(x -> {
for (int i = 0; i < 100000; i++) {
target.add(i);
}
});
System.out.println(target.size());
}
}
Результат определенного выполнения:163913. Если вы продолжите выполнять этоmainметод, в конечном итоге получит не-200000В итоге проблема здесь в использовании параллельных потоковparallelStream()метод.ParallelStreamдно используетсяFork/JoinРеализация фреймворка, то есть применяется пул потоковForkJoinPoolАбстрагируйте узлы в параллельном потоке какForkJoinTaskДля расчета принципы «кражи задач» и другие принципы, лежащие в их основе, здесь не будут распространяться, нужно только прояснить:
-
ForkJoinPoolОбщее использованиеRuntime.getRuntime().availableProcessors()(Это значение обычно считается количеством логического физического ядра машины) как степень параллелизма (parallelism), который просто считается количеством задач, которые могут выполняться одновременно, а не количеством рабочих потоков. - На многоядерных машинах используйте
ParallelStreamВсе операции в узле потока эквивалентны"в многопоточной среде"Выполнять операции, все операции в нем будут давать непредсказуемые результаты, такие как возможный выход массива за границы, недостающие добавленные элементы, частичные индексыindexцитируется какNULLи Т. Д.
Пример моделирования
Написание этой статьи не является преднамеренным, на самом деле автор уже давно столкнулся с относительно скрытым производственным сбоем, среди которых кусок кода с относительно низким трафиком выглядит примерно так:
@Data
private static class OrderDTO {
private String orderId;
private OrderStatus orderStatus;
private BigDecimal amount;
private Long customerId;
}
@Data
private static class Order {
private Long id;
private String orderId;
private Integer orderStatus;
private BigDecimal amount;
private Long customerId;
private OffsetDateTime createTime;
private OffsetDateTime editTime;
}
public void groupByOrderStatus(Long customerId) {
List<Order> orders = orderDao.selectByCustomerId(customerId);
List<OrderDTO> orderDTOList = new ArrayList<>();
orders.parallelStream().forEach(order -> {
OrderDTO dto = new OrderDTO();
......
orderDTOList.add(dto);
});
Map<String, List<OrderDTO>> collect
= orderDTOList.stream().collect(Collectors.groupingBy(item -> item.getOrderStatus().getCode()));
......
}
Функция метода передается через клиентаIDЗапросите список заказов, а затем преобразуйте список заказов вOrderDTOСписок, а затем сгруппирован по полю статуса заказа. Через журнал производства и регрессию по тестированию обнаружили, что вышеуказанный фрагмент кодаgroupByOrderStatus()Метод даже отправит аномалию пустого указателя.
Когда проблема впервые возникла, потому что разработчик передалLambdaВыражение сжимает несколько кодов в 1 строку, поэтому трудно проверить код с конкретной проблемой из стека исключений.LambdaВыражение разбивается на несколько строк в начальной точке периода.После наблюдения в течение определенного периода времени сегмент кода, в котором возникает исключение нулевого указателя, наконец находится.Collectors.groupingBy(item -> item.getOrderStatus().getCode()), это,OrderDTOв случаеorderStatusявляется пустым объектом. Здесь очевидно,groupByOrderStatus()На самом деле метод вызывается в стеке потоков, не должно быть нескольких потоков для одновременного изменения содержимого, здесь есть только одно сомнение: использованиеparallelStream(). Затем непосредственноparallelStream()превратиться вstream()Перезапустите, проблема с нулевым указателем больше не воспроизводится.
Lambda/StreamНа самом деле, это не естественно потокобезопасный.Предпосылкой потокобезопасности является то, что они закрываются и вызываются потоками и не вводят многопоточную среду.Например, если используются параллельные потоки, суть состоит в том, чтобы ввести многопоточная среда. Поэтому при разработке функций нужно хорошо подумать:
- Действительно ли необходимо использовать
Lambdaи потоковое программирование? - Действительно ли необходимо использовать параллельные потоки? Если используются параллельные потоки, нужно ли мне рассматривать возможность введения дополнительных механизмов синхронизации, таких как блокировки?
- Если вводится дополнительный механизм синхронизации, считается ли, что параллельный поток используется принудительно, что нарушает первоначальный замысел дизайна параллельного потока?
- На самом деле параллелизм не улучшает производительность, он может только улучшить пропускную способность Мы должны сосредоточиться на обнаружении и оптимизации узких мест производительности, а не отчаянно преобразовывать восходящий поток в параллельные вызовы.
❝Автор имеет привычку к чистоте кода.В то время я также обнаружил, что в приведенном выше коде есть операция сопоставления.Чтобы быть корректным, вместо forEach() следует использовать функцию map() для обхода элементов и перезагрузки их в другой список.Логика в методе отражает исходного разработчика.На самом деле, я мало что знаю о Lambda.
❞
резюме
Возвращаясь к исходному вопросу, на самом деле использование параллельных потоков также может обеспечить соответствие результатов выполнения ожиданиям, но необходимо ввести дополнительные механизмы синхронизации, такие как использование"монитор"Для синхронизации:
public class ParallelStreamMain {
public static void main(String[] args) throws Exception {
List<List<Integer>> array = new ArrayList<>();
List<Integer> item1 = new ArrayList<>();
List<Integer> item2 = new ArrayList<>();
List<Integer> target = new ArrayList<>(100);
array.add(item1);
array.add(item2);
final Object monitor = new Object();
array.parallelStream().forEach(x -> {
synchronized (monitor) {
for (int i = 0; i < 100000; i++) {
target.add(i);
}
}
});
System.out.println(target.size());
}
}
Независимо от того, сколько раз вышеуказанный метод будет выполняться, он выведет только:200000. Логика добавления синхронизированных блоков кода в параллельный поток здесь действительно кажется забавной, просто чтобы проиллюстрировать, что если элементы добавляются или изменяются в контейнере в многопоточной среде, только добавление дополнительного механизма синхронизации может гарантировать конечный результат. как и ожидалось. ParallelStream — это очень хороший дизайн, но вам нужно рассмотреть его применимые сценарии и не попасть в ловушку параллелизма, которую вы сами для себя расставили.