Сегодня я хочу поговорить о реактивных потоках. В последние годы этот подход становится всё популярнее. Spring позволяет использовать его вместо модели, основанной на сервлетах. Клиентские библиотеки также всё чаще предоставляют реактивный API, а базы данных поддерживают подключения через R2DBC и реактивные драйверы. Я тоже всё чаще применяю Project Reactor в своих модулях.
Реактивные потоки иногда сравнивают с Java Streams: обе технологии связаны с функциональным программированием и внешне похожи.
Однако обычный Stream можно потребить только один раз, поэтому разработчики редко возвращают из методов значения вроде Stream<String>.
Чаще весь конвейер обработки помещают в один метод. Для небольших преобразований данных это удобно, но такие вызовы в конечном итоге остаются блокирующими.
Несколько месяцев назад я решил пройти курс по Project Reactor на Udemy, чтобы систематизировать знания и найти новые идеи. Курс оказался интересным и полезным, особенно для разработчиков, которые только знакомятся с реактивным подходом.
После этого я подготовил реализацию «Игры Жизнь» на основе реактивных потоков. Это удачный пример: источник может генерировать бесконечную последовательность, а потребитель сам ограничивает количество элементов. Код проекта доступен в репозитории GameOfLife.
Проект состоит из нескольких частей:
CellState— перечисление, определяющее, жива клетка или мертва;Game— класс с основной логикой и генератором игровых циклов;UI— класс, выводящий поле в консоль.
Поле представлено двумерным массивом CellState[][].
Единственный публичный метод принимает начальное состояние и возвращает Flux последующих состояний, вычисляемых на его основе:
Flux<CellState[][]> game(CellState[][] initialField)
Метод намеренно не задаёт количество итераций — это решение принимает клиент. Поток строится в два этапа:
Mono.just(initialField)добавляет начальное поле в начало последовательности..concatWith(generate)присоединяет генератор следующих поколений.
Flux<CellState[][]> generate = Flux.generate(() -> initialField, (state, sink) -> {
CellState[][] iterate = iterate(state);
sink.next(iterate);
return iterate;
});
Генератор вычисляет новое поле из предыдущего, отправляет его в текущий Flux и сохраняет как состояние для следующего шага.
Тест показывает, как использовать эту логику:
@Test
public void printerTest() {
new Game().game(getGliderField())
.take(5)
.doOnNext(UI::printState)
.subscribe();
}
Игра запускается с начальной фигурой «планер», после чего поток ограничивается пятью элементами и каждое поле выводится в консоль.
Генератор выполнится четыре раза, потому что первым элементом уже является начальное состояние.
Если понадобится больше поколений или дополнительная обработка, новые операторы можно добавить непосредственно в существующий Flux.