Spring WebFlux — Spring Reactive Programming
Spring WebFlux — это новый модуль, представленный в Spring 5. Spring WebFlux — это первый шаг к модели реактивного программирования в среде Spring.
Spring реактивное программирование
Если вы новичок в модели реактивного программирования, я настоятельно рекомендую вам ознакомиться со следующими статьями, чтобы узнать о реактивном программировании.
- Реактивный манифест
- Реактивные потоки
- Реактивные потоки Java 9
- RxJava
Если вы новичок в Spring 5, ознакомьтесь с возможностями Spring 5.
Весенний веб-флюкс

- Mono: реализует Publisher и возвращает 0 или 1 элемент
- Flux: реализует Publisher и возвращает N элементов.
Пример Spring WebFlux Hello World

Зависимости Spring WebFlux Maven
4.0.0 com.journaldev.spring SpringWebflux 0.0.1-SNAPSHOT Spring WebFlux Spring WebFlux Example UTF-8 UTF-8 1.9 org.springframework.boot spring-boot-starter-parent 2.0.1.RELEASE org.springframework.boot spring-boot-starter-webflux org.springframework.boot spring-boot-starter-test test io.projectreactor reactor-test test spring-snapshots Spring Snapshots https://repo.spring.io/snapshot true spring-milestones Spring Milestones https://repo.spring.io/milestone false spring-snapshots Spring Snapshots https://repo.spring.io/snapshot true spring-milestones Spring Milestones https://repo.spring.io/milestone false org.springframework.boot spring-boot-maven-plugin org.apache.maven.plugins maven-compiler-plugin 3.7.0 $ $
Наиболее важными зависимостями являются spring-boot-starter-webflux и spring-boot-starter-parent . Некоторые другие зависимости предназначены для создания тестовых случаев JUnit.
Весенний обработчик WebFlux
Метод Spring WebFlux Handler обрабатывает запрос и возвращает Mono или Flux в качестве ответа.
package com.journaldev.spring.component; import org.springframework.http.MediaType; import org.springframework.stereotype.Component; import org.springframework.web.reactive.function.BodyInserters; import org.springframework.web.reactive.function.server.ServerRequest; import org.springframework.web.reactive.function.server.ServerResponse; import reactor.core.publisher.Mono; @Component public class HelloWorldHandler < public MonohelloWorld(ServerRequest request) < return ServerResponse.ok().contentType(MediaType.TEXT_PLAIN) .body(BodyInserters.fromObject("Hello World!")); >>
Обратите внимание, что реактивный компонент Mono содержит тело ServerResponse . Также посмотрите на цепочку функций, чтобы установить тип возвращаемого содержимого, код ответа и тело.
Весенний маршрутизатор WebFlux
Метод маршрутизатора используется для определения маршрутов для приложения. Эти методы возвращают объект RouterFunction , который также содержит тело ServerResponse .
package com.journaldev.spring.component; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.http.MediaType; import org.springframework.web.reactive.function.server.RequestPredicates; import org.springframework.web.reactive.function.server.RouterFunction; import org.springframework.web.reactive.function.server.RouterFunctions; import org.springframework.web.reactive.function.server.ServerResponse; @Configuration public class HelloWorldRouter < @Bean public RouterFunctionrouteHelloWorld(HelloWorldHandler helloWorldHandler) < return RouterFunctions.route(RequestPredicates.GET("/helloWorld") .and(RequestPredicates.accept(MediaType.TEXT_PLAIN)), helloWorldHandler::helloWorld); >>
Итак, мы предоставляем метод GET для /helloWorld , и клиентский вызов должен принимать простой текстовый ответ.
Весеннее загрузочное приложение
Давайте настроим наше простое приложение WebFlux с помощью Spring Boot.
package com.journaldev.spring; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; @SpringBootApplication public class Application < public static void main(String[] args) < SpringApplication.run(Application.class, args); >>
Если вы посмотрите на приведенный выше код, там нет ничего, связанного с Spring WebFlux. Но Spring Boot настроит наше приложение как Spring WebFlux, поскольку мы добавили зависимость от модуля spring-boot-starter-webflux .
Поддержка модулей Java 9
Наше приложение готово к выполнению на Java 8, но если вы используете Java 9, нам также необходимо добавить класс module-info.java .
module com.journaldev.spring
Запуск приложения Spring WebFlux Spring Boot

2018-05-07 15:01:47.893 INFO 25158 --- [ main] o.s.w.r.f.s.s.RouterFunctionMapping : Mapped ((GET && /helloWorld) && Accept: [text/plain]) -> com.journaldev.spring.component.HelloWorldRouter$$Lambda$501/704766954@6eeb5d56 2018-05-07 15:01:48.495 INFO 25158 --- [ctor-http-nio-1] r.ipc.netty.tcp.BlockingNettyContext : Started HttpServer on /0:0:0:0:0:0:0:0:8080 2018-05-07 15:01:48.495 INFO 25158 --- [ main] o.s.b.web.embedded.netty.NettyWebServer : Netty started on port(s): 8080 2018-05-07 15:01:48.501 INFO 25158 --- [ main] com.journaldev.spring.Application : Started Application in 1.86 seconds (JVM running for 5.542)
Из логов видно, что наше приложение работает на сервере Netty на порту 8080. Давайте продолжим и протестируем наше приложение.
Весенний тест приложения WebFlux
Мы можем протестировать наше приложение различными методами.

Резюме
В этом посте мы узнали о Spring WebFlux и о том, как создать реактивный веб-сервис Hello World Restful. Приятно видеть, что популярные фреймворки, такие как Spring, поддерживают модель реактивного программирования. Но нам нужно охватить многое, потому что, если все ваши зависимости не являются реактивными и неблокирующими, ваше приложение также не является по-настоящему реактивным. Например, поставщики реляционных баз данных не имеют реактивных драйверов, потому что они зависят от JDBC, который не является реактивным. Следовательно, Hibernate API также не является реактивным. Поэтому, если вы используете реляционные базы данных, вы пока не можете создать действительно реактивное приложение. Я надеюсь, что это изменится раньше, чем позже.
Вы можете скачать код проекта из моего репозитория GitHub.
Ссылка: Официальная документация
Java Blog
Spring WebFlux — это новый реактивный веб-фреймворк, представленный в Spring Framework 5.0. В отличие от Spring MVC, он не требует API сервлета, является полностью асинхронным и неблокирующим, а также реализует спецификацию Reactive Streams через проект Reactor.
Spring WebFlux поставляется в двух вариантах: функциональном и на основе аннотаций. Основанная на аннотациях модель довольно близка к модели Spring MVC, как показано в следующем примере:
@RestController @RequestMapping("/users") public class MyRestController < @GetMapping("/") public Mono getUser(@PathVariable Long user) < // . >@GetMapping("//customers") public Flux getUserCustomers(@PathVariable Long user) < // . >@DeleteMapping("/") public Mono deleteUser(@PathVariable Long user) < // . >>
“WebFlux.fn”, функциональный вариант, отделяет конфигурацию маршрутизации от фактической обработки запросов, как показано в следующем примере:
@Configuration(proxyBeanMethods = false) public class RoutingConfiguration < @Bean public RouterFunctionmonoRouterFunction(UserHandler userHandler) < return route(GET("/").and(accept(APPLICATION_JSON)), userHandler::getUser) .andRoute(GET("//customers").and(accept(APPLICATION_JSON)), userHandler::getUserCustomers) .andRoute(DELETE("/").and(accept(APPLICATION_JSON)), userHandler::deleteUser); > > @Component public class UserHandler < public MonogetUser(ServerRequest request) < // . >public Mono getUserCustomers(ServerRequest request) < // . >public Mono deleteUser(ServerRequest request) < // . >>
WebFlux является частью Spring Framework.
Вы можете определить столько bean-компонентов RouterFunction, сколько захотите, чтобы сделать определение маршрутизатора модульным. Beans можно упорядочить, если вам нужно применить приоритет.
Для начала добавьте модуль spring-boot-starter-webflux в ваше приложение.
Добавление в приложение модулей spring-boot-starter-web и spring-boot-starter-webflux приводит к автоматической настройке Spring Boot Spring MVC, а не к WebFlux. Это поведение было выбрано потому, что многие разработчики Spring добавляют spring-boot-starter-webflux в свое приложение Spring MVC для использования реактивного WebClient. Вы можете применить свой выбор, установив для выбранного типа приложения значение SpringApplication.setWebApplicationType(WebApplicationType.REACTIVE)
- Spring Boot: разработка веб-приложений
- Разработка вашего первого Spring Boot приложения
- Spring Boot стартеры
Реактивное программирование на Java. Будущее, настоящее и прошлое
Разберемся с парадигмой реактивного программирования. Какие есть плюсы и минусы по сравнению с императивным подходом.
28 мая 2023 · 14 минуты на чтение
В современном мире разработки, где требования к производительности и отзывчивости системы становятся все более строгими, реактивное программирование выделяется как ключевой подход.
Что же такое реактивное программирование, и почему оно становится всё более популярным? В чём недостаток императивного подхода? И, что самое важное, как реактивные потоки помогают нам создавать более производительные и эффективные системы?
Целью статьи является погружение в контекст реактивной разработки и объяснение основных механик. Это не гайд по написанию реактивных приложений. Это будет в следующих статьях.
Спонсор поста
Реактивная система
Так как реактивный подход помогает создавать реактивные системы, неплохо сначала разобраться, что это за системы такие.
Представим себе систему управления таксопарка. Владелец таксопарка идёт в ногу со временем и решил заказать разработку системы для управления заказами. Система включает в себя множество подсистем: работа с заказами, управлением автопарком и так далее.
Перед разработкой владелец подсчитал среднее количество заказов в день, аналитики по этим данным рассчитали необходимое количество железа, заложили в эти данные избыток в 50% на будущий рост и закупили сервера. Система была написана и введена в эксплуатацию.
Всё было хорошо, пока в городе не объявили проведение чемпионата мира по футболу. Толпы туристов, многие из которых решили воспользоваться удобным способом заказа такси. В какой-то момент нагрузка превзошла все самые смелые ожидания и система полностью развалилась. В итоге таксопарк потерял клиентов и прибыль, а его рейтинг в AppStore обвалился.
Кажется, что было бы неплохо, чтобы система как-то динамически реагировала на изменения, будь то резко выросшая нагрузка или недоступность внешних служб.
В этом примере стоит задуматься об увеличении эластичности. Пропускная способность вашей системы должна автоматически увеличиваться при увеличении нагрузки, и автоматически уменьшаться со снижением нагрузки.
Это одно из свойств, которыми должна обладать реактивная система. Поговорим обо всех характеристиках.
Отзывчивость
Система способна быстро обрабатывать запросы пользователей даже при высокой нагрузке. Это требует соблюдения нескольких ключевых принципов проектирования.
Неблокирующий ввод/вывод: Использование неблокирующего ввода-вывода, позволит минимизировать время, которое потоки тратят на ожидание завершения операций ввода-вывода. Более эффективное использование потоков снижает вероятность «голодания» потоков и увеличивает производительность сервиса.
Про неблокирующий ввод/вывод я расскажу ниже более подробно.
Балансировка нагрузки гарантирует, что ни один экземпляр не будет перегружен трафиком. Таким образом увеличивая пропускную способность и производительность системы.
Существуют различные алгоритмы распределения запросов, но их цель одна: распределить запросы равномерно на существующие и работоспособные экземпляры системы.
Кэширование: система должна использовать методы кэширования для сокращения времени, затрачиваемого на обработку запросов, и повышения общей производительности системы.
Допустим, у вас есть сервис справочников, который содержит значения для различных списков: статусы заказа, названия категорий товаров. Нет смысла каждый раз обращаться в сервис справочной информации, если данные там меняются нечасто. Кэширование этой информации позволит снизить затраты на межсервисное взаимодействие, а также позволит сервису продолжить работу, даже если сервис справочников будет недоступен.
Для больших систем отличным решением будет использование CDN для кэширования статического контента и снижения сетевых задержек.
Устойчивость
Система должна продолжать работать во время сбоев и автоматически восстанавливаться после ошибок, а не полностью выходить из строя.
Существует несколько ключевых стратегий, которые могут помочь обеспечить устойчивость:
Отказоустойчивость: Реактивная система должна быть спроектирована таким образом, чтобы выдерживать сбои на всех уровнях, включая аппаратные, сетевые и программные сбои. Это может быть достигнуто благодаря избыточности, репликации и механизмов автоматического распределения нагрузки с отказавших экземпляров сервисов на здоровые экземпляры.
Автоматическое восстановление: Реактивная система должна уметь автоматически обнаруживать и диагностировать ошибки, а также предпринимать корректирующие действия для восстановления после сбоев без вмешательства разработчиков.
Вынужденная деградация (Graceful degradation): Система должна продолжать работать, даже если некоторые её компоненты или функции недоступны или работают неправильно.
Вместо того чтобы полностью разрушиться, система плавно деградирует, отключая несущественные функции, снижая функциональность или предоставляя запасные варианты. Это позволяет системе продолжать работать, хотя и с ограниченными возможностями, пока проблема не будет устранена.
Предохранители (Circuit breakers): Модель проектирования, которая может помочь предотвратить каскадные отказы в системе. Работает путём мониторинга количества отказов, происходящих в сервисе за определённый период времени, и автоматически отключает предохранитель, если количество отказов превышает заданный порог.
Можно реализовать различные предохранители. Например, если сервис не отвечает на запрос, начать возвращать дефолтное значение или кэшированное значение. Все зависит от ваших сценариев.
Контроль потока данных (Backpressure): Механизм, позволяющий получателю управлять скоростью получения данных от отправителя. Иными словами, это метод контроля потока информации между отправителем и получателем.
Включив эти стратегии в конструкцию реактивной системы, можно достичь высокого уровня устойчивости и гарантировать, что система способна быстро восстанавливаться после ошибок и продолжать бесперебойную работу.
Эластичность
Реактивная система должна иметь возможность масштабирования для обработки растущих рабочих нагрузок без снижения производительности и доступности.
Существует несколько методов, которые могут быть использованы для достижения эластичности в реактивной системе:
Горизонтальное масштабирование: Предполагает добавление дополнительных экземпляров сервиса для распределения рабочей нагрузки на несколько машин или узлов.
Этот подход может использоваться для обработки растущего трафика или запросов пользователей и может быть достигнут при использовании технологий контейнеризации, таких как Docker, и инструментов оркестрации, таких как Kubernetes.
Вертикальное масштабирование: Предполагает добавление дополнительных ресурсов (процессор, память или дисковое пространство) к существующему сервису для увеличения его мощности.
Этот подход может быть использован для обработки возросших объёмов данных или требований к обработке, и может быть реализован с помощью облачных инфраструктурных сервисов, таких как Amazon EC2 или Microsoft Azure.
Автомасштабирование: Предполагает использование автоматизированных инструментов или алгоритмов для динамической корректировки количества экземпляров или ресурсов, выделяемых сервису, на основе показателей, собираемых в реальном времени, таких как использование процессора, памяти или сетевого трафика.
Динамическая корректировка может увеличить количество экземпляров сервиса при увеличении количества запросов, а также уменьшить количество экземпляров при уменьшении нагрузки. Увеличение позволяет обработать непредсказуемый рост нагрузки, а уменьшение позволяет эффективнее потреблять доступные ресурсы и не платить накладные расхода за не используемые.
Достижение эластичности в реактивной системе требует сочетания архитектурного проектирования, управления инфраструктурой и инструментов автоматического масштабирования для обеспечения того, чтобы система могла адаптироваться к изменяющимся рабочим нагрузкам и поддерживать свою производительность и доступность в течение долгого времени.
Управление сообщениями
Message Driven — шаблон проектирования, используемый в реактивных системах для обеспечения асинхронной связи между различными компонентами. Этот паттерн позволяет сервисам отправлять и получать сообщения без блокировки, ожидания или удержания ресурсов, что помогает минимизировать потребление ресурсов и максимизировать пропускную способность.
Для обеспечения коммуникации на основе сообщений реактивные системы обычно используют брокер сообщений или очередь сообщений, которая выступает в качестве посредника между различными сервисами. Когда сервис посылает сообщение другому сервису, он помещает его в очередь сообщений, а принимающий компонент может получить сообщение асинхронно, когда он будет готов обработать его.
Такой подход позволяет сервисам работать независимо друг от друга, без необходимости знать о состоянии или доступности других сервисов. Он также позволяет сервисам обрабатывать большие объёмы сообщений или событий, не перегружаясь и не блокируясь входящим трафиком.
Некоторые популярные брокеры сообщений и системы очередей, используемые в реактивных системах: Apache Kafka, RabbitMQ и AWS SQS.
Подписывайся на Telegram
Реактивное программирование
Манифест реактивных систем гласит: «Большие системы состоят из подсистем, имеющих те же свойства и, следовательно, зависят от их реактивных характеристик. Это означает, что принципы Реактивных Систем применяются на всех уровнях.» Таким образом, каждый отдельный сервис должен также следовать принципам реактивной системы.
Реактивное программирование — это парадигма программирования, ориентированная на работу с потоками данных и распространение изменений в этих потоках. Какие-то данные поступают в систему, и как реакция на них, система выполняет какие-то действия. Вот отсюда и название — «реактивное».
Хотя реактивное программирование часто используется для создания реактивных систем, технически возможно достичь определённой степени реактивности системы, используя другие методы программирования, такие как многопоточность, асинхронное программирование и событийно-ориентированные архитектуры. Эти методы могут помочь улучшить отзывчивость и масштабируемость системы, но они могут не обеспечить тот уровень устойчивости и отказоустойчивости, который призваны обеспечить реактивные системы.
Проблемы императивного программирования
Разберёмся в недостатках императивного подхода. Зачем понадобилось выдумывать какое-то реактивное программирование, почему сложно написать реактивную систему на существующих технологиях?
Реализуем небольшой пример, который состоит из двух сервисов.
@Service @RequiredArgsConstructor public class PassengerServiceImpl implements PassengerService < private final RideService rideService; @Override public void requestRide(Location pickupLocation, Location dropoffLocation) < RideRequest rideRequest = new RideRequest(pickupLocation, dropoffLocation); Ride ride = rideService.processRideRequest(rideRequest); // other logic >>
PassengerServiceImpl представляет API для запроса поездки, ориентированный на пассажира. Метод requestRide() принимает в качестве параметров место посадки и высадки пассажира и инициирует процесс запроса поездки.
@Service @RequiredArgsConstructor public class RideServiceImpl implements RideService < private final DriverService driverService; @Override public Ride processRideRequest(RideRequest rideRequest) < Driver driver = driverService.findAvailableDriver(rideRequest.getPickupLocation()); // other logic Ride ride = new . return ride; >>
RideServiceImpl представляет внутреннюю службу, которая обрабатывает запросы на поездки. Метод processRideRequest() принимает объект RideRequest в качестве параметра и инициирует процесс поиска водителя и назначения поездки.
Представим, что DriverService при вызове метода findAvailableDriver() обращается к базе данных или посылает сетевой запрос в другой сервис. Что будет, если БД будет выполнять запрос 30 секунд или другой сервис ответит спустя 5 минут?
Одна из основных проблем императивного подхода это ожидание потоков выполнения какой-либо задачи, то есть блокировка. Например, для выполнения запроса к БД из пула потоков берётся поток, далее он ожидает , пока БД выполнит запрос и вернёт результат. Если вычисление результата займёт 5 минут, то поток всё это время будет недоступен для других операций.
Это может привести к снижению производительности сервиса, особенно если многие потоки будут блокироваться в ожидании завершения долго выполняющихся запросов к базе данных. В какой-то момент у вас просто могут закончиться потоки в пуле, и обработка новых запросов просто остановится.
Такая же проблема может возникнуть, когда выполняется запрос к внешнему сервису. Например, вы посылаете запрос используя RestTemplate . Если внешний ресурс будет отвечать 5 минут, то всё это время поток будет находится в ожидании ответа, то есть простаивать.

Почему простаивание потока — это проблема?
Каждый поток нуждается в памяти для хранения своего стека вызовов и других связанных с ним структур данных. Когда поток простаивает, он продолжает потреблять ресурсы для поддержания своего состояния.
Кроме того, процессорное время, которое выделяется неработающим потокам, могло бы быть использовано для других задач. Если большое количество потоков простаивает, это может привести к увеличению загрузки процессора и снижению производительности, так как операционная система будет тратить больше времени на переключение между потоками.
Попробуем решить проблему блокировки потоков доступными способами. Добавим ExecutorService в RideService и будем возвращать не Ride , а Future .
public inteface RideService < public FutureprocessRideRequest(RideRequest rideRequest); >
@Service public class PassengerServiceImpl implements PassengerService < private final RideService rideService; @Override public void requestRide(Location pickupLocation, Location dropoffLocation) < RideRequest rideRequest = new RideRequest(pickupLocation, dropoffLocation); Futurefuture = rideService.processRideRequest(rideRequest); // other logic Ride ride = future.get(); // other logic > >
Теперь мы выполняем асинхронный вызов к RideService и получаем объект Future . Далее мы можем продолжить выполнять другие операции, пока выполняется обработка Future .
Мы можем выполнить какую-то другую логику, но в какой-то момент необходимо вызывать метод Future.get() , который потенциально также является блокирующим, если Future ещё не закончил работу, то мы заблокируем поток.
Эта реализация позволила нам сократить время блокировки потока, однако полностью эта проблема не решена, мы всё ещё с большой вероятностью будем получать блокировку потока.
Более высокоуровневым решением может быть использование CompletionStage и его реализации CompletableFuture . CompletionStage позволяет писать код в функциональном стиле, который выполняется асинхронно.
public inteface RideService < public CompletionStageprocessRideRequest(RideRequest rideRequest); >
@Service public class PassengerServiceImpl implements PassengerService < private final RideService rideService; public PassengerServiceImpl(RideService rideService) < this.rideService = rideService; >@Override public void requestRide(Location pickupLocation, Location dropoffLocation) < RideRequest rideRequest = new RideRequest(pickupLocation, dropoffLocation); rideService.processRideRequest(rideRequest) .thenApply(a ->< . >) .thenCombine(b -> < . >) .thenAccept(c -> < . >) > >
Однако, реализации с использованием Future , CompletionStage и им подобных требуют от разработчика глубокого понимания многопоточного программирования: доступ к общей памяти, синхронизация, обработка ошибок и так далее.
Но и это ещё не всё. Дизайн многопоточности в Java не предполагает, что мы будем создавать поток на каждый чих. Создание потока дорогостоящая операция. Да, пул потоков частично решает эту проблему, но есть ещё одна проблема: несколько потоков могут использовать один процессор для выполнения задач одновременно. При такой ситуации, процессорное время распределяется между несколькими потоками, что вызывает необходимость переключения контекста. Для возобновления выполнения потока позже, необходимо сохранять и загружать регистры, карты памяти и выполнить другие операции с высоким объёмом вычислений. Из-за этого снижается эффект от использования большого количества потоков при небольшом количестве процессоров.
Паттерн Наблюдатель (Observer Pattern)
Вспомним паттерн «Наблюдатель». Он поможет нам лучше понять концепцию реактивных потоков.
Наблюдатель — это поведенческий паттерн проектирования, который создаёт механизм подписки, позволяющий одним объектам следить и реагировать на события, происходящие в других объектах.
В этом паттерне есть два ключевых участника: Издатель и Подписчик (Наблюдатель). Издатель обновляет состояние и оповещает всех своих подписчиков об этих изменениях. Подписчики, в свою очередь, реагируют на эти уведомления.
Ключевые характеристики
- Декаплинг: Субъекты и наблюдатели функционируют независимо друг от друга. Это означает, что они не должны знать друг о друге. Субъекты просто отправляют уведомления, а наблюдатели просто реагируют на них.
- Динамичность: Подписчики могут подписываться и отписываться от субъектов в любое время.
- Многопоточность: Паттерн Наблюдатель позволяет обрабатывать события асинхронно и в различных потоках исполнения.
В контексте реактивного программирования паттерн Наблюдатель становится основой для создания реактивных потоков, обеспечивающих эффективную обработку данных и событий. Это также служит фундаментом для различных библиотек и фреймворков, таких как RxJava и Project Reactor.
Реактивные потоки (стримы)
Спецификация Reactive Streams впервые была опубликована в 2015 году. Она была разработана для стандартизации модели асинхронного потокового программирования с контролем потока данных в JVM и была принята в качестве основы для обработки асинхронного потока в JDK 9.
В отличие от обычных Java Stream не было предоставлено стандартных реализаций реактивных стримов, поэтому в последующие годы в Java-сообществе появилось несколько библиотек и фреймворков, которые реализуют и расширяют спецификацию Reactive Streams, таких как: RxJava, Vert.x, ProjectReactor, Akka Streams.
Базовая схема работы стримов
Рекомендую посмотреть доклад «Олег Докука — Реактивный хардкор: как построить свой Publisher», который закрепит понимание данного механизма взаимодействия.
Продолжая историю с паттерном Наблюдатель, теперь у нас есть интерфейсы Publisher , представляет источник данных, и Subscriber — наблюдатель, который подписывается на поток и получает уведомления об изменении состояния потока. Есть ещё один важный интерфейс, который является «посредником» между первыми двумя — это Subscription .
Давайте взглянем, как выглядят данные интерфейсы:
package org.reactivestreams; public interface Publisher < public void subscribe(Subscribers); >
package org.reactivestreams; public interface Subscriber
package org.reactivestreams; public interface Subscription
Все начинается с подписки Subscriber на Publisher посредством вызова метода Publisher.subscribe() . Publisher использует переданный объект Subscriber , вызывая метод onSubscribe() , передавая в Subscriber объект Subscription . Через этот объект Subscriber будет взаимодействовать с Publisher .

Теперь Subscriber , используя полученный Subscription , будет запрашивать значения у Publisher . Это важный момент, не Publisher инициализирует отправку данных подписчикам когда хочет, это подписчики запрашивают необходимое количество данных у Publisher . Таким образом реализуется контроль потока данных (Backpressure).

Метод onNext(T t) вызывается, когда Subcriber запрашивает значения у Publisher , используя метод Subscription.request() . Он передаёт данные по одному, но не больше, чем было запрошено подписчиком.
Метод onError() вызывается, когда ошибка происходит на стороне Publisher , оповещая таким образом о проблеме Subscriber , передавая объект исключения.
Метод onComplete() оповещает подписчиков, что у Publisher не осталось элементов для передачи. Данный метод, как и onError() вызывается лишь один раз.
Пример реактивного потока
Представим, что у нас есть набор чисел, и мы хотим получить квадрат каждого числа. При этом числа для расчёта могут поступать из разных систем или из бд, или из БД редиса и других систем.
В императивном стиле программирования мы бы обработали эти данные следующим образом:
List numbers = externalServices.getNumbers(); for (int number : numbers)
Код обрабатывает данные в определённом порядке и от начала до конца. Мы должны дождаться пока метод getNumbers() отправит нам все данные для расчёта.
Давайте рассмотрим тот же пример с использованием Project Reactor:
Flux numbers = externalServices.getNumber(); // Flux это реализация Publisher numbers .map(number -> number * number) .subscribe( data -> System.out.println("Получены данные: " + data), error -> System.out.println("Произошла ошибка: " + error), () -> System.out.println("Поток данных завершен") );
Этот код напоминает Stream API. Только вместо Stream используется Flux — это тип данных из Project Reactor, который представляет собой поток данных.
Мы используем метод map , чтобы преобразовать каждое число в его квадрат, и затем подписываемся на поток, чтобы вывести результат. Подписка очень важна, без неё ничего не произойдёт. Только в отличие от Stream API на Flux можно подписаться несколько раз.
В данном примере мы будем обрабатывать данные асинхронно по мере их поступления. Одна система ответила быстрее другой, сразу обработали эту часть данных.
Ограничения и недостатки
Не бывает идеального решения и реактивные потоки не исключение.
Чтобы получить преимущество реактивных потоков, весь стек должен быть реактивным: доступ к БД, операции чтения/записи файлов и так далее. Всё должно работать в реактивной парадигме, иначе вы получите блокировки, которые ухудшат производительность всей системы.
Например, стандартный JDBC не является реактивным. Если использовать его в реактивном сервисе, то придётся ждать ответ, когда мы отправляем запрос в базу данных. Соответственно, вся реактивность тут же ломается.
Весь технологический стек должен быть реактивным
Сложность: Реактивное программирование может быть сложным для понимания, особенно для новичков. Это связано с необходимостью работы с асинхронностью, обработкой ошибок и механизмами подобными backpressure.
Отладка и тестирование: Отладка и тестирование реактивных систем сложная задача, так как асинхронная природа реактивного кода делает отслеживание исполнения программы менее прямолинейным и понятным.
Требования к проектированию: Реактивные системы требуют продуманного проектирования, чтобы обеспечить их эффективность. Без правильного проектирования система может столкнуться с проблемами производительности или недостаточной отзывчивости.
Поддержи автора
Event Loop
Рассуждая на тему реактивного программирования, нельзя пройти мимо такого понятия, как Event Loop.
Это реактивная асинхронная модель программирования для серверов. Она позволяет достичь более высокого уровня параллелизма при меньшем количестве потоков.
По сути, Event Loop — это реализация шаблона Reactor. Является неблокирующим потоком ввода-вывода, который работает непрерывно. Его основная задача — проверка новых событий. И как только событие пришло перенаправлять его тому, кто в данный момент может его обработать. Иногда их может быть несколько для увеличения производительности.

Выше приведён абстрактный дизайн цикла событий, который представляет идеи реактивного асинхронного программирования:
- Цикл событий выполняется непрерывно в одном потоке, хотя у нас может быть столько циклов событий, сколько доступно ядер.
- Цикл событий последовательно обрабатывает события из очереди событий и возвращается сразу после регистрации обратного вызова в платформе.
- Платформа может инициировать завершение операции, такой как вызов базы данных или вызов внешней службы.
- Цикл событий может запускать обратный вызов при уведомлении о завершении операции и отправлять результат обратно исходному вызывающему.
В своей работе механизм Event Loop использует Netty — клиент-серверная среда ввода-вывода для разработки сетевых приложений Java. Также этот механизм использует Vert.x – это полифункциональная библиотека для построения реактивных приложений на JVM.
Реактивные фреймворки в Java
Теперь немного поговорим про имплементации реактивной спецификации, коих уже появилось достаточное количество.
Spring WebFlux (Project Reactor)
Project Reactor – это библиотека для реактивного программирования на Java, которая полностью поддерживает Reactive Streams. Она предлагает два основных типа данных – Flux и Mono .
Flux представляет поток ноль или более элементов, а Mono представляет один или ноль элементов. Оба типа предоставляют обширный набор операторов для трансформации и комбинирования этих потоков.
Spring WebFlux – это веб-фреймворк, который является частью экосистемы Spring и использует Project Reactor для обработки реактивных потоков. В отличие от Spring MVC, который предназначен для синхронного веб-программирования и блокирования I/O, Spring WebFlux предназначен для асинхронного и неблокирующего веб-программирования.
Вместо Tomcat используется неблокирующий Netty. Сервлеты тоже ушли в прошлое.

Quarkus (Vert.x)
Не спрингом единым. Активно использую его в работе и в пет-проектах.
Реактивное программирование в Quarkus основано на библиотеке SmallRye Mutiny. Mutiny предлагает два основных типа: Uni и Multi , которые представляют собой типы реактивных потоков, обеспечивающих обработку одного или множества элементов соответственно.
Множество знакомых коллег, которым довелось поработать и со Spring WebFlux и с Quarkus, отмечают, что реактивный Quarkus API более приятный для работы.
Quarkus также оптимизирован для работы в контейнеризированных средах и облачных приложениях, что, в сочетании с его реактивной моделью программирования, делает его отличным выбором для создания современных микросервисов.
Кроме того, Quarkus предоставляет возможность компиляции приложений в нативный код при помощи GraalVM, что позволяет уменьшить время запуска и использование памяти, что является ещё одним преимуществом при использовании в облачных средах.
Project Loom
Некоторые разработчики предрекают скорую смерть таким проектам, как Project Reactor, ведь уже близится релиз Project Loom. Давайте порассуждаем на эту тему.
Цель Project Loom добавить в Java так называемые виртуальные потоки или нити (fibers). Одной из ключевых особенностей виртуальных потоков является их «непрерывность». Виртуальный поток может быть приостановлен, а его ресурсы могут быть возвращены в пул потоков, что позволяет использовать поток операционной системы для другой работы. Когда придёт ответ от внешнего сервиса, виртуальный поток может быть возобновлен и продолжить свою работу.
Хотя такие реактивные библиотеки предоставляют мощные абстракции для управления асинхронным и неблокирующим кодом, они всё ещё полагаются на традиционную модель потоков Java, которая может быть сложной и трудной для понимания. С Project Loom разработчики получат более простой способ написания неблокирующего кода, без необходимости использования сложных абстракций или пулов потоков.
Однако, на данный момент Project Loom находится в разработке, а вот реактивные фреймворки уже есть и успешно используются.
Вот аргументы, почему существующие фреймворки останутся в строю:
- Project Loom находится в разработке, и его окончательная форма и влияние на экосистему Java не полностью понятны.
- Ничто не мешает реактивным фреймворкам использовать под капотом новые виртуальные потоки. Предоставляя мощные абстракции и API для работы.
- Когда требуется тип обработки в стиле событий (Event-Driven Architecture), то их API очень удобен, странно от него отказываться.
Рекомендую посмотреть следующие доклады на тему Project Loom:
- Иван Углянский — Thread Wars: проект Loom наносит ответный удар
- Олег Докука, Андрей Родионов — Project Loom — друг или враг Reactive?
Заключение
Мы разобрались, что такое реактивные системы, и какими свойствами система должна обладать, чтобы называться реактивной. Разработка реактивной системы — сложный процесс и очевидно, что не каждой системе необходимо быть реактивной.
Однако, системы написанные с применением реактивных подходов и реактивного стека позволяют выдерживать большую нагрузку, чем системы, написанные на стандартном императивном стеке, а также быть более эффективными с точки зрения потребления ресурсов системы и дальнейшего её масштабирования.
Создать реактивную систему проще всего с помощью реактивных фремворков, которые позволяют не блокировать потоки и обрабатывать данные по мере их поступления.
Дополнительные материалы
- Зачем нам Reactive и как его готовить. Команда разработки делится своим опытом перевода сервисов на реактивный стек. Они используют Spring WebFlux.
Reactive Programming: Reactor и Spring WebFlux — часть 1
Я очень запарился по теме реактивного программирования, что наконец решился написать эту статью 🙂
Очень надеюсь она будет полезна для кого-то.Когда я начал разбираться с этой темой у меня было много вопросов на этот счет: что это такое, зачем ?
Прочитав некоторые статьи на эту тему я не получил полноту знаний, чтобы просто понять, что такое реактивная система.
Я решил углубиться в данную тему: перечитал все возможные туториалы и статьи на этот счет и прочел официальную документацию.
Сейчас попробую рассказать своими (и не только) словами, как же я понимаю, что такое реактивная система.
Я постарался изложить максимально детально некоторые очевидные вещи, но по ходу статьи можно будет найти много чего интересного.Надеюсь после это статьи у читателя не останется лишних вопросов.
Большую часть материала из этой статьи можно найти на просторах интернета и в официальной документации.
Я постарался максимально сжато и в то же время подробно остановиться на главных и ключевых моментах.
Но это вовсе не значит, что можно остановиться на прочтении данной статьи.
Впереди еще много всего интересного, что касается реактивщины 🙂
Кто еще такой этот Reactor ?
Project Reactor — это библиотека Java 8, которая реализует модель реактивного программирования. Он построен на основе спецификации реактивных потоков , стандарта для создания реактивных приложений.
Reactor — это полностью неблокирующая основа реактивного программирования для JVM с эффективным управлением требованиями (в форме управления «противодавлением»). Он напрямую интегрируется с функциональными API-интерфейсами Java 8, в частности, Completable Future , Stream и Duration . Он предлагает составные API-интерфейсы асинхронной последовательности — Flux (для элементов [N]) и Mono (для элементов [1]) — и широко реализует спецификацию Reactive Streams.
Reactor — это реализация парадигмы реактивного программирования, которую можно описать следующим образом:
Wiki:
Реактивное программирование — это парадигма асинхронного программирования, связанная с потоками данных и распространением изменений. Это означает, что становится возможным легко выражать статические (например, массивы) или динамические (например, излучатели событий) потоки данных через используемый язык (языки) программирования.
Ядро реактора работает на Java 8 и выше.
Все примеры и реализации Mono и Flux, касающиеся Reactor и WebFlux, которые мы рассматриваем, будут одинаковыми. Далее мы разберемся со всем в контексте WebFlux, который построен на Reactor.
Деталь проектов:
- Ядро Reactor — Реактивные основы для приложений и сред, а также реактивные расширения, основанные на API с типами Mono (1 элемент) и Flux (n элементов).
- Reactor Netty — предлагает неблокирующие и готовые к backpressure TCP / HTTP / UDP клиенты и серверы на основе Netty фреймворка.
- Reactor Addons — Мост к RxJava 2 Observable, Completable, Flowable, Single, Maybe, Scheduler, а также SWT Scheduler, Akka Scheduler и так далее.
Что такое реактивные типы и зачем их использовать ?
Реактивные типы не предназначены для обработки запросов или данных быстрее.
Их сила заключается в их способности одновременно обслуживать больше запросов и более эффективно обрабатывать операции с задержкой, такие как запрос данных с удаленного сервера.
Они позволяют обеспечить лучшее качество обслуживания и предсказуемое планирование пропускной способности, изначально имея дело со временем и задержкой, не затрачивая больше ресурсов.
В отличие от традиционной обработки, которая блокирует текущий поток во время ожидания результата, Reactive API запрашивает только тот объем данных, который он способен обработать и предоставляет новые возможности, поскольку он имеет дело с потоком данных, а не с отдельными элементами (объектами).
Следуя документации Spring WebFlux, Spring Framework использует Reactor для своей собственной реактивной поддержки.
Сейчас мы более детально во всем разберемся.
Reactor — это реализация Reactive Streams, которая дополнительно расширяет базовый контракт Reactive Streams Publisher с типами API-интерфейсов, которые можно комбинировать с Flux и Mono.
Reactive Streams (Реактивные Потоки)
Reactive Streams (Реактивные Потоки) состоят из 4-х простых Java — интерфейсов (Publisher, Subscriber, Subscription и Processor).
public static interface PublisherT> public void subscribe(Subscribersuper T> subscriber);
> public static interface SubscriberT> public void onSubscribe(Subscription subscription);
public void onNext(T item);
public void onError(Throwable throwable);
public void onComplete();
>
public static interface Subscription public void request(long n);
public void cancel();
>
public static interface ProcessorT,R> extends SubscriberT>,
PublisherR> >
К ним всем выдвигаются примерно следующие требования:
ASYNC — асинхронность
NIO — “неблокируемость” ввода/вывода
RESPECT BACKPRESSURE — умение обрабатывать ситуации, когда данные появляются быстрее, чем потребляются (в синхронном, императивном коде подобная ситуация не возникает, но в реактивных системах такое часто встречается).
Различия с Java 8
Давайте немного вспомним.
До java 8 были Future и Callbacks.
Concurrency API ввел понятие сервиса-исполнителя (ExecutorService). Исполнители выполняют задачи асинхронно и обычно используют пул потоков, так что нам не надо создавать их вручную.
Кроме Runnable , исполнители могут принимать другой вид задач, который называется Callable. Но как тогда получить результат, который они возвращают? Поскольку метод submit() не ждет завершения задачи, исполнитель не может вернуть результат задачи напрямую. Вместо этого исполнитель возвращает специальный объект Future, у которого мы сможем запросить результат задачи.
Callback методы являются void (ничего не возвращают) и принимают дополнительный параметр, который вызывается после определенного события.
Интерфейс Future описывает API для работы с задачами, результат которых мы планируем получить в будущем. Для Future нас интересует его реализация java.util.concurrent.FutureTask. То есть это Task, который будет выполнен во Future. Чем эта реализация ещё интересна, так это тем, что она реализует и Runnable. Можно считать это своего рода адаптером старой модели работы с задачами в потоках и новой модели (появилась в java 5).
ExecutorService executor = . ;
Future f = executor.submit(. );
f.get();
Но дело в том, что Future является асинхронным, но блокирует текущий поток до тех пор, пока не будет завершено вычисление при попытке получить результат с помощью метода get ().
Future API был хорошим шагом на пути к асинхронному программированию, но ему не хватало некоторых важных и полезных функций.
- Future нельзя завершить вручную.
Допустим, у вас есть метод для получения свободных номеров в гостинице из удалённого API. Поскольку этот вызов API занимает много времени, вы запускаете его в отдельном потоке и возвращаете Future.
Теперь предположим, что удалённый сервис перестал работать и вы хотите завершить Future вручную, передав актуальную цену продукта из кэша. К сожалению вы не сможете этого сделать.
- Нельзя выполнять дальнейшие действия над результатом Future без блокирования.
Также в Future нельзя повесить функцию-колбэк, чтобы она срабатывала автоматически, как только станет доступен результат.
- Невозможно выполнить множество Future один за другим.
Такой алгоритм асинхронной работы невозможен при использовании Future.
- Невозможно объединить несколько Future.
Именно по этому Spring Framework 4 представил ListenableFuture — это Future реализация, которая добавляет неблокирующие возможности на основе обратного вызова.
У Guava есть интерфейс ListenableFuture, который является будущим, к которому вы можете прикрепить обратный вызов, который будет вызываться, когда будет доступен результат, так что вам не придется вызывать get() и создавать блок потока, пока результат не станет доступен.
С ListenableFuture вы можете зарегистрировать callback так:
ListenableFuture listenable = service.submit(. );
Futures.addCallback(listenable, new FutureCallback() @Override
public void onSuccess(Object o) //handle on success
>
@Override
public void onFailure(Throwable throwable) //handle on failure
>
>)
Технически ListenableFuture расширяет интерфейс Future, добавляя простые:
void addListener (Runnable listener, Executor executor);
Потом в java 8 появились CompletableFuture и лямбды и облегчили жизнь многим разработчикам.
CompletableFuture позволяет иметь дело с будущим неблокирующим образом обеспечивая возможности для цепочки отложенной обработки результатов.
CompletableFuture используется для асинхронного программирования в Java. Асинхронное программирование — это средство написания неблокирующего кода путём выполнения задачи в отдельном, отличном от главного, потоке, а также уведомление главного потока о ходе выполнения, завершении или сбое.
Таким образом, основной поток не блокируется и не ждёт завершения задачи, а значит может параллельно выполнять и другие задания.
CompletableFuture реализует интерфейс Future, то есть наши task будут выполнены в будущем, и мы сможем выполнить get() и получить результат. Но ещё он реализует CompletionStage.
CompletableFuture запускает цепочку на выполнение сразу, не дожидаясь того, что у него попросят посчитанное значение, в отличие от Stream API, где при создании стрима он не запускается сразу, а ждёт, когда из него захотят значение.
Надо помнить, что CompletalbeFuture в своей работе использует Runnable, Consumer и Function.
Класс CompletableFuture является строго реализацией CompletionStage, имеющей дело с асинхронными вычислениями, которые завершатся и предоставят значение в некоторое неуказанное время.
С CompletableFuture вы также можете зарегистрировать callback, когда задача завершена, но она отличается от ListenableFuture тем, что она может быть завершена из любого потока, который хочет ее выполнить:
CompletableFuture completableFuture = new CompletableFuture();
completableFuture.whenComplete(new BiConsumer() @Override
public void accept(Object o, Object o2) //handle complete
>
>); // complete the task
completableFuture.complete(new Object())
Для того, чтобы асинхронно выполнить некоторую фоновую задачу, которая не возвращает, результат, можно использовать метод CompletableFuture.runAsync(). Он принимает объект Runnable и возвращает CompletableFuture.
CompletableFuture future = CompletableFuture.runAsync(() -> try TimeUnit.SECONDS.sleep(100);
> catch (InterruptedException e) throw new IllegalStateException(e);
>
System.out.println("Work is done in a separate thread");
>);
Для выполнения асинхронной задачи и возврата результата стоит использовать CompletableFuture.supplyAsync(). Он принимает Supplier и возвращает CompletableFuture, где T это тип возвращаемого функцией-поставщиком значения:
CompletableFuture future = CompletableFuture.supplyAsync(() -> try TimeUnit.SECONDS.sleep(100);
> catch (InterruptedException e) throw new IllegalStateException(e);
>
return "Result of an Asynchronous Task";
>);String result = future.get();
System.out.println(result);
runAsync() и supplyAsync() выполняются в отдельном потоке.
CompletableFuture выполняет эти задачи в потоке, полученном из глобального ForkJoinPool.commonPool().
Вы можете повесить callback на CompletableFuture, используя методы thenApply(), thenAccept() и thenRun().
- thenApply() служит для обработки и преобразования результата CompletableFuture при его поступлении. В качестве аргумента он принимает Function.
CompletableFuture future= CompletableFuture.supplyAsync(() -> try TimeUnit.SECONDS.sleep(100);
> catch (InterruptedException e) throw new IllegalStateException(e);
>
return "Hello";
>);
CompletableFuture greetingFuture = future.thenApply(name -> return name + "World!";
>);
System.out.println(greetingFuture.get()); // Hello, World!
Если вы не хотите возвращать результат, а хотите просто выполнить часть кода после завершения Future, можете воспользоваться методами thenAccept() и thenRun(). Эти методы являются потребителями и часто используются в качестве завершающего метода в цепочке.
CompletableFuture.supplyAsync(() -> return StudentService.getCurrentStudents(studentId);
>).thenAccept(student -> System.out.println("Get all info about current srudent" + student.getFirstName())
>);
Мы также можем объеденить несколько CompletableFuture, обрабатывать исключения, использовать асинхронные callback и другие возможности.
Но это все не предназначено для работы с задержкой, такой как операции ввода-вывода. И именно здесь появляются такие Reactive API, как Reactor или RxJava.
Реактивные API, такие как Reactor, предназначены для обработки как синхронных, так и асинхронных операций и позволяют буферизовать, объединять или применять широкий спектр преобразований к вашим данным.
Изначально API Reactive были разработаны только для работы с потоками данных типа Flux.
Но со временем был представлен и поток Mono.
Mono является реактивным эквивалентом CompletableFuture типа.
Чуть дальше мы более детально разберем Mono и Flux.
Flux и Mono реализуют Publisher интерфейс из спецификации Reactive Streams.
Оба класса соответствуют спецификации, и мы могли бы использовать этот интерфейс вместо них:
Publisher just = Mono.just("one");
Но на самом деле это было сделано потому что некоторые операции имеют смысл только для одного из двух типов.
Основной задачей Reactive Streams является обработка backpressure. Я не стал переводить это слово, дабы не допустить ошибки в понимании.
Backpressure — это механизм, который позволяет получателю спрашивать, сколько данных он хочет получить.
Т.е. получатель начинает получать данные только тогда, когда он готов их обработать.
Основным артефактом проекта Reactor является reactor-core реактивная библиотека, которая фокусируется на спецификации Reactive Streams и ориентирована на Java 8+
Принцип работы
В официальной документации Reactor сравнивается с конвейером.
Publisher выдаёт какие-то данные (материалы). Данные идут по цепочке из операторов (конвейерной ленте), обрабатываются, в конце получается готовый продукт, который передаётся в нужный Consumer/Subscriber и употребляется уже там.
Оператор — это некий Publisher, который помимо какой-то своей логики содержит ссылку на исходный Publisher, к которому применяется. Вызовы операторов создают цепочку из Publisher.
Реактивное программирование возникло из-за желания писать асинхронный неблокируемый код в читаемом виде. Ни код, написанный на колбэках, ни код, написанный с помощью CompletableFuture, не может быть настолько удобочитаемым, как этого можно добиться при помощи реактивности.
В основе подхода лежит идея разделения компонентов на 2 типа: источник событий (Publisher) и обработчик событий (Subscriber).
Subscriber подписывается на события, которые создаёт Publisher, а затем каким-то образом их обрабатывает. По сути это паттерн Observer с надстроенными поверх возможностями и особенностями.
Существует еще одно понятие — Observer.
Он может подписаться на событие объекта и выполнять какие либо действия с полученным результатом.
У одного Subject может быть много подписчиков.
Общение между Publisher и Subscriber происходит через объект Subscription.
Subcriber может регулировать скорость поставки сообщений от Publisher (т.к. backpressure), а также отменять подписку.
Publisher’ы можно объединять в цепочки и комбинировать разными способами.
Интерфейсы Publisher, Subscriber и Subscription находятся в пакете org.reactivestreams, который по умолчанию был добавлен в Java 9.
Они задают спецификацию для собственной реализации реактивных потоков.
Библиотека Project Reactor является такой реализацией.
Небольшое вступление
В Reactor есть два типа Publisher: Flux и Mono.
О них чуть ниже.
Как только происходит публикация ивента, Subscriber его сразу же и получает.
Как можно получить поток ?
Разными способами, например
Flux flux1 = Flux.just(“foo”, “bar”, “foobar”);
Flux flux2 = Flux.fromIterable(Arrays.asList(“A”, “B”, “C”));
Flux flux3 = Flux.range(5, 3);
На самом деле, при вызове нового каждого оператора в цепочке создаётся новый Publisher, который добавляется к цепочке.
Например:
Flux.just(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
.limitRequest(5)
.skip(3)
.subscribe(value -> System.out.println("Value: " + value));
В качестве аналогии можно привести конвейер, по которому перемещаются элементы. Операторы же по очереди проводят с ними манипуляции.
Вначале мы ограничили лимит = 5, пропустили 3 элемента и подписались на наш поток, чтобы получить данные.
Пока все просто.
Следующие элементы в цепочке вызываются под капотом при помощи метода onNext().
Что касается исключений, так они являются терминальными. Это значит что поток сразу завершает свою работу.
Вы можете вернуть значение по умолчанию в случае ошибки выполнения метода при помощи onErrorReturn().
Никаких дополнительных потоков по умолчанию не создаётся при использовании реактивного подхода. Всё работает в потоке подписки: в каком потоке мы подписываемся на Publisher, в том все элементы и будут обрабатываться.
Спецификация реактивных потоков запрещает null значения в последовательности.
Ну а теперь разберемся чуть более подробно.
Что такое WebFlux ?
Spring 5 представил платформу WebFlux, которая представляет собой полностью асинхронный и неблокирующий реактивный веб-стек, который позволяет обрабатывать огромное количество одновременных соединений.
Модуль WebFlux является альтернативой Spring MVC и представляет собой реактивный подход для написания веб-сервисов.
WebFlux позиционирует себя как микрофреймворк.
Этот новый микрофреймворк поддерживает аннотированные контроллеры, функциональные конечные точки, WebClient (аналог RestTemplate в Spring Web MVC), WebSockets и многое другое.
В основе WebFlux лежит библиотека Reactor.
Вам нужен Spring Boot версии 2+ для использования модуля Spring WebFlux.
WebFlux по умолчанию использует Netty. Tomcat не поддерживает реактивщину.
Пару слов о Netty
Он предназначен для неблокирующего ввода-вывода.
На входе у него есть один поток, который работает в бесконечном цикле. Благодаря селектору и канальному механизму, он перенаправляет данные из входящих запросов во входящие буферы и делегирует обработку этих запросов выделенному пулу потоков асинхронных потоков. И также в обратном направлении.
В чем отличие Spring WebFlux от RxJava ?
Spring WebFlux и RxJava2+ являются реализациями реактивных потоков.
Т.е. теперь начинает складываться более детальная картина:
Project Reactor — это библиотека, которая лежит в основе WebFlux.
Она исправляет недостатки в RxJava и больше подходит для бэкэнд-разработки. RxJava имеет некоторые проблемы, которые могут вызвать нехватку памяти, например (взято из официальных источников).
Как сказал David Karnok
Use Reactor 3 if you are allowed to use Java 8+, use RxJava 2 if you are stuck on Java 6+ or need your functions to throw checked exceptions.
Если можно сказать простыми словами, то “Реактивное программирование” — это программирование с асинхронными потоками (streams) данных.
Т.е. можно слушать поток и реагировать на события в нем. Можем фильтровать поток как нам вздумается, объединять потоки, кроме того потоки могут быть входными параметрами других потоков.
Даже множественный поток может быть использован как входной аргумент другого потока. Вы можете объединять несколько потоков. Вы можете фильтровать один поток, чтобы потом получить другой, который содержит только актуальные данные. Вы можете объединять данные с одного потока с данными другого, чтобы получить еще один.
Поток — это некая последовательность, состоящая из постоянных событий, отсортированных по времени. В нем может быть три типа сообщений: значения (данные некоторого типа), ошибки и сигнал о завершении работы.
Нам лишь надо подписаться на поток.
Наблюдатель (observers) подписывается на поток.
Данные же будут получены тогда, когда они будут готовы — произойти это может в этом же самом либо другом потоке. Publisher их сам отдаст, когда они придут. Это push-модель. Отдача готового элемента — это event. Реактивная модель основана на событиях (event-driven).
Что такое Back-pressure ?
Мы возвращаем не объект, а “обещание” объекта — Publisher, который будет отдавать объекты, как только они появятся. Отдавать мы их будем Subscriber-у — тому, кто подписывается на Publisher.
Подписчик может быть как один, так и много. Subscriber и получает объекты. Причем подписчик может регулировать скорость потока, это и называется Back-pressure.
Как уже говорилось выше, основная концепция реактивного программирования — это неблокирующий ввод/вывод.
Слушать поток означает подписаться на него. Т.е. функции, которые мы определили — это наблюдатели (Observers). А поток является субъектом, который наблюдают. Этот подход называется Observer Design Pattern.
Если у нас есть Publisher, который отправляет события потребителю быстрее, чем он может их обработать, то, в конце концов, потребитель будет перегружен событиями, которые истощают системные ресурсы. Backpressure означает, что наш клиент должен иметь возможность сообщить производителю, сколько данных отправлять, чтобы предотвратить это, и это указано в самой спецификации.
Попробуем реализовать этот механизм следующим образом.
Давайте скажем восходящему потоку отправлять только два элемента одновременно, используя request()
List elements = new ArrayList<>();
Flux.just(1, 2, 3, 4)
.log()
.subscribe(new Subscriber() private Subscription s;
int onNextAmount;
@Override
public void onSubscribe(Subscription s) this.s = s;
s.request(2);
>
@Override
public void onNext(Integer integer) elements.add(integer);
onNextAmount++;
if (onNextAmount % 2 == 0) s.request(2);
>
>
@Override
public void onError(Throwable t) <>
@Override
public void onComplete() <>
>);
16:26:39.915 [main] INFO reactor.Flux.Array.1 — | onSubscribe([Synchronous Fuseable] FluxArray.ArraySubscription)
16:26:39.917 [main] INFO reactor.Flux.Array.1 — | request(2)
16:26:39.917 [main] INFO reactor.Flux.Array.1 — | onNext(1)
16:26:39.917 [main] INFO reactor.Flux.Array.1 — | onNext(2)
16:26:39.917 [main] INFO reactor.Flux.Array.1 — | request(2)
16:26:39.917 [main] INFO reactor.Flux.Array.1 — | onNext(3)
16:26:39.917 [main] INFO reactor.Flux.Array.1 — | onNext(4)
16:26:39.917 [main] INFO reactor.Flux.Array.1 — | request(2)
16:26:39.918 [main] INFO reactor.Flux.Array.1 — | onComplete()
По сути, это реактивное baclpressure. Мы просим стрим подтолкнуть нам только определенное количество элементов, и только тогда, когда когда мы будем готовы.
Зачем мне вообще использовать RxJava или Reactor, если то же самое можно сделать с Streams, CompletableFutures и Optionals?
Проблема, по сути, заключается в том, что большую часть времени вы имеете дело с простыми задачами и вам действительно не нужны эти библиотеки. Но когда все усложняется, вы должны писать некрасивый кусок кода. Затем этот кусок кода становится все более сложным и сложным в поддержке.
RxJava и Reactor имеют множество удобных функций, которые будут удовлетворять ваши потребности на долгие годы.
Есть два варианта получить данные
— Pull
Это когда мы сами делаем запрос на получение и нам приходит ответ.
— Push
Когда данные сами нас уведомляют об изменениях и система “выталкивает” их нам.
Реактивное приложение, это когда приложение само извещает нас об изменении своего состояния. Не мы делаем запрос и проверяем, а не изменилось ли там что-то, а приложение само нам сигнализирует. Ну и конечно эти события и эти сигналы мы можем обрабатывать как нам вздумается.
Pull — коллекция (аналог — массив): в ней есть данные, которые мы можем получить по запросу, предварительно обработав их как нам хочется.
Push — полная противоположность: изначально в ней нет данных, но как только они появятся, она нам сообщит об этом. Во время этого мы также можем делать с ней что хотим и как только в коллекции появятся значения, она выполнит все наши фильтры (которые мы на нее навесили) и выдаст нам результат.
Push коллекция как-бы “состоит” из новой сущности Observable.
Это и есть коллекция, которая будет рассылать уведомления об изменении своего состояния.
Для реализации этого похода в реактивном контексте существует такое понятие как Callback — это объект или несколько объектов, которые “отслеживают” необходимые события, происходящие с обслуживаемыми объектами, и либо сообщают об этих событиях другим слушателям, либо отдают объекты, с которыми произошли это события, асинхронным потоком для обработки.
Спецификация для реактивного подхода Spring WebFlux:
Основные концепции
В новом подходе у нас есть два основных класса для работы в реактивном режиме:
— Mono
Класс Mono нужен для работы с единственным объектом.
Mono также может использоваться как какая-то асинхронная задача в стиле “выполнил и забыл”, без возвращаемого результата (очень похож на Runnable).
Данный класс схож с Mono, но предоставляет возможность асинхронной работы со множеством объектов.
Flux — это Publisher, способный выпустить от 0 до N событий (элементов), в том числе и бесконечное их число.
Flux и Mono реализуют Publisher интерфейс из спецификации Reactive Streams.
Разделение на Flux и Mono помогает улучшить семантику реактивного API, делая его достаточно выразительным.
Flux и Mono — lazy.
Для того чтобы запустить какую-то обработку и воспользоваться данными, лежащими в Mono и Flux, нужно на них подписаться с помощью subscribe().
Методы subscribe() используют “лямбда-выражения” из Java 8 в качестве параметров.
Способы подписаться
- Подписаться (слушать поток):
subscribe();
2) Сделать что-то с каждым полученным значением
subscribe(Consumer consumer);
3) Сделать что-то в случае исключения
subscribe(Consumer consumer, Consumer errorConsumer);
4) Сделать что-то по завершению
subscribe(
Consumer consumer,
Consumer errorConsumer,
Runnable completeConsumer
);
Данные передаются следующим образом
- Вызывается метод subscribe()
2) Затем создается Subscription объект
3) После этого Subscriber вызовет request() метод в Subscription классе, чтобы указать количество объектов, которые он может обработать (если этот метод не вызывается явно, запрашивается неограниченное количество объектов).
4) Subscriber может получать объекты с помощью методаonNext().
Если подписчик получает все запрошенные им объекты, он может запросить дополнительные объекты или отменить подписку, вызвав onComplete(). Если в какой-то момент возникает ошибка, издатель вызывает метод onError() на подписчике.
Flux.just(1, 2, 3, 4)
.subscribe(System.out::println);
Данные не начнут поступать, пока мы не подпишемся — метод subscribe().
Когда мы вызываем подписку, под капотом вызывается Subscription, который запрашивает элементы из потока (это означает, что он запрашивает каждый доступный элемент).
Этот поток описывается в интерфейсе Subscriber как часть спецификации реактивных потоков, и фактически это то, что было реализовано за кулисами в нашем вызове метода onSubscribe().
В интерфейсе Subscriber 4 метода:
— onSubscribe()
— onNext()
— onComplete()
— onError()
Мы можем записать это по-другому
Flux.just(1, 2, 3, 4)
.subscribe(new Subscriber() @Override
public void onSubscribe(Subscription s) s.request(Long.MAX__VALUE);
>
@Override
public void onNext(Integer integer) elements.add(integer);
>
@Override
public void onError(Throwable t) <>
@Override
public void onComplete() <>
>);
Давайте глянем на следующий пример
List elements = new ArrayList<>();
Flux.just(1, 2, 3, 4)
.log()
.subscribe(elements::add);
Мы получим следующий результат
20:25:19.550 [main] INFO reactor.Flux.Array.1 — | onSubscribe([Synchronous Fuseable] FluxArray.ArraySubscription)
20:25:19.553 [main] INFO reactor.Flux.Array.1 — | request(unbounded)
20:25:19.553 [main] INFO reactor.Flux.Array.1 — | onNext(1)
20:25:19.553 [main] INFO reactor.Flux.Array.1 — | onNext(2)
20:25:19.553 [main] INFO reactor.Flux.Array.1 — | onNext(3)
20:25:19.553 [main] INFO reactor.Flux.Array.1 — | onNext(4)
20:25:19.553 [main] INFO reactor.Flux.Array.1 — | onComplete()
onSubscribe() —вызывается когда мы подписываемся на поток
request(unbounded) — вызывается когда мы вызываем подписку (subscribe), за кулисами мы создаем Subscription. Эта подписка запрашивает элементы из потока. В этом случае он запрашивает каждый доступный элемент.
onNext() — вызывается на каждом элементе
onComplete() — вызывается последним, после получения последнего элемента.
Основные приемущества
Реактивность дает слабую связанность.
В некоторых случаях это дает возможность писать более простой и понятный код.
Например мы можем взять обычную коллекцию, преобразовать ее к реактивной коллекции и тогда мы будем иметь коллекцию событий об изменении данных в ней. Мы очень просто получаем только те данные, которые изменились. По этой коллекции мы можем делать выборку, фильтровать и т.д.
Если бы мы это делали обычным традиционным способом, то нам нужно было бы получить данные, закэшировать их, сделать запрос на получение новых данных,сравнить их с текущим значением.
Чтобы получить данные, мы обращаемся к базе данных, получаем данные и отдаем пользователю. Если у вас в этот момент времени на секунду пропадет интернет то вы получите ошибку. Беда.
Реактивщина нам поможет в данном случае.
Аналог для примера: Gmail, Facebook, Instagram, Twitter. Когда у вас плохой интернет, вы не получаете ошибку, а просто ждете результат немного дольше.
Вы едете в метро, обновили ленту в instagram, у у вас пропал интернет, спустя минуту он у вас появился и вам не надо будет обновлять пальцем сверху вниз (сделать еще один запрос на получение данных), система вам отдаст эти данные сама.
Обязанности:
— Приложение должно отдавать пользователю результат за полсекунды
— Обеспечить отзывчивость под нагрузкой
— Система остается в рабочем состоянии даже, если один из компонентов отказал.
— Система должна занимать оптимальное количество ресурсов в каждый промежуток времени.
Общение между сервисами должно происходить через асинхронные сообщения. Это значит, что каждый элемент системы запрашивает информацию из другого элемента, но не ожидает получение результата сразу же. Вместо этого он продолжает выполнять свои задачи.
Когда стоит использовать ?
Реактивщину стоит использовать, когда есть поток событий, растянутый во времени (например, пользовательский ввод).
Функционал, написанный при помощи реактивных потоков, может быть легко дополнен и расширен.
Но самое важное, это то, что не надо использовать реактивщину везде где только попало)
Реактивщина + Реляционные БД
Многие в интернете говорят, что транзакционная БД не подходит для реактивной концепции. Концепция транзакции не совсем соответствует реактивному миру, так как она связана с блокировкой ресурса, а это именно то, чего стараются избежать, используя реактивность.
Например если вы используете реляционную БД без поставляемого реактивного драйвера, то все запросы к этой базе всё равно будут блокируемыми, так что никакой выгоды от использования реактивного стека к этим запросам здесь не будет, как бы этого не хотелось.
На данный момент предпринимаются первые шаги в реализации реактивного подхода для реляционных баз данных, хотя готовых решений ещё не существует.
Реактивный клиент для SQL DB:
ссылка
Некоторые примеры с Mono и Flux
Я хотел бы привести несколько примеров для полноты понимания работы потоков Mono и Flux.
После того, как Flux/Mono отправил какие-то данные, подписка на эти события происходит просто: вызываем метод subscribe().
Также мы можем передать туда лямбду, которая будет вызываться для каждого элемента в потоке
Flux.just("1", "2", "3", "4")
.subscribe(value -> System.out.println("Value: " + value));
В случае ошибки вызовется метод onError().
Этот обработчик можно задать вторым параметром subscribe() метода
Flux.just(1, 2, 3, 4, 5)
.subscribe(value -> if (value > 4) throw new IllegalArgumentException(value + " > than 4");
>
System.out.println("Value: " + value);
>, error -> System.out.println("Error: " + error.getMessage()));илиFluxInteger> ints = Flux.range(1, 4)
.map(i -> if (i 3) return i;
throw new RuntimeException("Got to 4");
>);
ints.subscribe(System.out::println,
error -> System.err.println("Error: " + error));
В случае успешного завершения выполнится обработчик onComplete().
Его можно задать третьим параметром
Flux.just(1, 2, 3, 4)
.subscribe(value -> System.out.println("Value: " + value),
error -> <>,
() -> System.out.println("Successfull"));
Четвертый вариант работы с subscribe() методом, это Consumer . Этот вариант требует, чтобы вы что-то сделали с Subscription (выполнить request(long) на нем или cancel()), иначе Flux просто повиснет. Определяем его как четвертый параметр.
Здесь мы говорим, что мы хотим до 10 элементов из источника (который на самом деле испустит 4 элемента и завершится).
Flux ints = Flux.range(1, 4);
ints.subscribe(System.out::println,
error -> System.err.println("Error " + error),
() -> System.out.println("Done"),
sub -> sub.request(10));
Экземпляры типа Mono и Flux можно “преобразовывать” друг в друга.
Например Flux.collectList() вернет Mono, а Mono.flux() — вернет Flux.
А метод block() — блокируем пока не будет получен следующий сигнал или не истечет время ожидания (метод перегружен: public T block(Duration timeout))
Mono> listMono = Flux.just(1, 2, 3, 4)
.filter(value -> value % 2 == 0)
.collectList();System.out.println(listMono.block());
Вы можете отменить подписку с помощью объекта Disposable.
Метод onDispose() может использоваться для отмены подписки после завершения Flux или ошибок или очистки.
В результате получите: 1, 2
Disposable disposable = Flux.just(1, 2, 3, 4, 5, 6, 7, 8)
.delayElements(Duration.ofSeconds(3))
.subscribe(value -> System.out.println("Value: " + value));
Thread.sleep(7000);
disposable.dispose();
System.out.println("Cancelling subscription");
These variants [of operators] return a reference to the subscription that you can use to cancel the subscription when no more data is needed. Upon cancellation, the source should stop producing values and clean up any resources it created. This cancel and clean-up behavior is represented in Reactor by the general-purpose Disposable interface.
Вы также можете отменить подписку с помощью onCancel().
Метод onCancel() может использоваться для выполнения любых действий, специфичных для отмены до того как работает onDispose().
Метод onCancel() вызывается первым — можно использовать для выполнения любых действий, относящихся к отмене до работы onDispose().
Метод onDispose() можно использовать для выполнения очистки, когда Flux завершает работу, выдает ошибки или отменяется.
Вы можете создать свой Flux
— при помощи метода .generate()
— при помощи метода .create()
— при помощи just()
— justOrEmpty — выводит Mono<>, если элемент не null, иначе сигнал завершения
— fromArray
— fromIterable
— range
— fromStream
— fromSupplier
— fromRunnable
— fromFuture
— из стороннего Publisher
FluxString> seq1 = Flux.just("foo", "bar", "foobar");
ListString> iterable = Arrays.asList("foo", "bar", "foobar");
FluxString> seq2 = Flux.fromIterable(iterable);MonoString> noData = Mono.empty(); // пустой MonoMonoString> data = Mono.just("foo");FluxInteger> numbersFromFiveToSeven = Flux.range(5, 3); // Первый параметр - это начало диапазона, а второй параметр - количество элементов, которые нужно произвести.
FluxInteger> ints = Flux.range(1, 3);
ints.subscribe(); // подписка на поток
или
FluxInteger> ints = Flux.range(1, 3);
ints.subscribe(i -> System.out.println(i)); // будет выводить значения в консоль. Получим: 1, 2, 3
PublisherString> publisher = redisson.getKeys().getKeys();
FluxString> from = Flux.from(publisher);
Что такое Reactor IPC ?
Reactor IPC — это расширение Reactor, которое позволяет интегрироваться с различными платформами и системами, поддерживающими неблокирующий ввод / вывод (например, Netty, Kafka).
Когда запрос поступает в Netty, он немедленно обрабатывается с использованием класса ChannelOperations . Далее вызывается цепочка вызовов и, наконец, запрос достигает DispatcherHandler (он направляет входящие запросы другим контроллерам).
Затем запрос достигает контроллера.
Дело в том, что по умолчанию все Mono и Flux холодные (ниже мы рассмотрим более подробно, что это означает).
Затем на основе Publisher начал выстраиваться поток, достигший класса ChannelOperations, и только здесь, в ChannelOperations, был вызван метод subscribe(), и только в этот момент этот поток начал работать (начали происходить HTTP-вызовы).
Спасибо что дочитали до конца.
Если вы нашли неточности в описании данной статьи, вы можете написать мне на email и я с радостью вам отвечу.
Kirill Sereda