Павло Васильченко

Павло Васильченко

Principal Member of Technical Staff

Русский

Начало работы с реактивным программированием

Published August 4, 2022

Сегодня я хочу поговорить о реактивных потоках. В последние годы этот подход становится всё популярнее. Spring позволяет использовать его вместо модели, основанной на сервлетах. Клиентские библиотеки также всё чаще предоставляют реактивный API, а базы данных поддерживают подключения через R2DBC и реактивные драйверы. Я тоже всё чаще применяю Project Reactor в своих модулях.

Реактивные потоки иногда сравнивают с Java Streams: обе технологии связаны с функциональным программированием и внешне похожи. Однако обычный Stream можно потребить только один раз, поэтому разработчики редко возвращают из методов значения вроде Stream<String>. Чаще весь конвейер обработки помещают в один метод. Для небольших преобразований данных это удобно, но такие вызовы в конечном итоге остаются блокирующими.

Несколько месяцев назад я решил пройти курс по Project Reactor на Udemy, чтобы систематизировать знания и найти новые идеи. Курс оказался интересным и полезным, особенно для разработчиков, которые только знакомятся с реактивным подходом.

После этого я подготовил реализацию «Игры Жизнь» на основе реактивных потоков. Это удачный пример: источник может генерировать бесконечную последовательность, а потребитель сам ограничивает количество элементов. Код проекта доступен в репозитории GameOfLife.

Проект состоит из нескольких частей:

  • CellState — перечисление, определяющее, жива клетка или мертва;
  • Game — класс с основной логикой и генератором игровых циклов;
  • UI — класс, выводящий поле в консоль.

Поле представлено двумерным массивом CellState[][]. Единственный публичный метод принимает начальное состояние и возвращает Flux последующих состояний, вычисляемых на его основе:

Flux<CellState[][]> game(CellState[][] initialField)

Метод намеренно не задаёт количество итераций — это решение принимает клиент. Поток строится в два этапа:

  1. Mono.just(initialField) добавляет начальное поле в начало последовательности.
  2. .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.