Пользовательская реализация основных концепций реактивного программирования на Java (аналог RxJava).
В проекте реализована система реактивных потоков с возможностью управления потоками выполнения и обработки событий, построенная на паттерне «Наблюдатель» (Observer). Реализованы базовые компоненты, операторы преобразования данных, планировщики потоков и механизмы отмены подписки.
-
RxObservable — источник данных, фабрики
create(),just(). -
RxObserver — интерфейс с методами
onNext(),onError(),onComplete(). -
Операторы (в пакете
com.rx.operators):MapOperator(map)FilterOperator(filter)FlatMapOperator(flatMap)MergeOperator(merge)ConcatOperator(concat)ReduceOperator(reduce)
-
Schedulers (в пакете
com.rx.schedulers):RxIOScheduler(cached thread pool)RxComputationScheduler(fixed thread pool)RxSingleScheduler(single-thread executor)
-
Disposable:
RxDisposable— отмена одной подпискиRxCompositeDisposable— групповая отмена
-
Логирование через SLF4J + Log4j
- Java 17+
- Maven
- SLF4J API + Log4j
- JUnit 5
-
Клонировать репозиторий:
-
Собрать и запустить тесты:
mvn clean test -
Запустить демонстрацию:
mvn exec:java -Dexec.mainClass="com.rx.Main"
-
Паттерн Observer:
- Источник (
RxObservable) делегирует эмиссию элементов черезRxOnSubscribe. - Потребитель реализует
RxObserverили передаёт лямбды вsubscribe(). RxDisposableконтролирует отмену,RxCompositeDisposable— групповую отмену.
- Источник (
-
Структура пакетов:
core— базовые компоненты и фабрики.operators— классы-операторы для модульности.schedulers— управление планировщиками потоков.
-
Flow:
- Построение цепочки:
RxObservable.create(...)→ операторы →subscribeOn()/observeOn()→subscribe(). - Все переходы потоков выполняются через
RxScheduler.schedule(...).
- Построение цепочки:
| Scheduler | Реализация | Применение |
|---|---|---|
| RxIOScheduler | CachedThreadPool |
I/O задачи, сеть |
| RxComputationScheduler | FixedThreadPool(N=CPU) |
CPU-bound вычисления |
| RxSingleScheduler | SingleThreadExecutor |
Последовательная обработка |
subscribeOn()определяет поток подписки.observeOn()переключает поток обработки событий.
В проекте написаны юнит-тесты JUnit 5 для ключевых сценариев:
-
Базовая работа
create()+subscribe(onNext, onError, onComplete)just(), проверка эмиссии и завершения.
-
Операторы
map,filterflatMap,merge,concat,reduce
-
Планировщики
subscribeOn/observeOnпроверяют переключение потоков.
-
Обработка ошибок
- Эмит
onError, проверка прекращенияonNext.
- Эмит
-
Отмена подписки
RxDisposable.dispose(),RxCompositeDisposable.dispose().
Запуск:
mvn test// map + filter + планировщики
MapOperator.apply(
RxObservable.just(1,2,3,4,5),
i -> i * 2
)
.subscribeOn(new RxIOScheduler())
.observeOn(new RxSingleScheduler())
.subscribe(
i -> System.out.println("-> " + i),
Throwable::printStackTrace,
() -> System.out.println("Done")
);
// flatMap
FlatMapOperator.apply(
RxObservable.just("A","B"),
s -> RxObservable.just(s + "1", s + "2")
).subscribe(System.out::println);