From b52e90861a657bbd23706c89441fe9911934775d Mon Sep 17 00:00:00 2001 From: Wosin Date: Tue, 20 Mar 2018 17:08:07 +0100 Subject: [PATCH] RxRelay Article Materials (#3859) --- rxjava/pom.xml | 5 + .../java/com/baeldung/rxjava/RandomRelay.java | 32 +++++ .../java/com/baeldung/rxjava/RxRelayTest.java | 120 ++++++++++++++++++ 3 files changed, 157 insertions(+) create mode 100644 rxjava/src/main/java/com/baeldung/rxjava/RandomRelay.java create mode 100644 rxjava/src/test/java/com/baeldung/rxjava/RxRelayTest.java diff --git a/rxjava/pom.xml b/rxjava/pom.xml index d88dfcaa9b..9a07aba2a3 100644 --- a/rxjava/pom.xml +++ b/rxjava/pom.xml @@ -71,6 +71,11 @@ 22.0 test + + com.jakewharton.rxrelay2 + rxrelay + 2.0.0 + diff --git a/rxjava/src/main/java/com/baeldung/rxjava/RandomRelay.java b/rxjava/src/main/java/com/baeldung/rxjava/RandomRelay.java new file mode 100644 index 0000000000..1367a23847 --- /dev/null +++ b/rxjava/src/main/java/com/baeldung/rxjava/RandomRelay.java @@ -0,0 +1,32 @@ +package com.baeldung.rxjava; + +import com.jakewharton.rxrelay2.Relay; +import io.reactivex.Observer; +import io.reactivex.disposables.Disposables; + +import java.util.ArrayList; +import java.util.List; +import java.util.Random; + +public class RandomRelay extends Relay { + Random random = new Random(); + + List> observers = new ArrayList<>(); + + @Override + public void accept(Integer integer) { + int observerIndex = random.nextInt(observers.size()) & Integer.MAX_VALUE; + observers.get(observerIndex).onNext(integer); + } + + @Override + public boolean hasObservers() { + return observers.isEmpty(); + } + + @Override + protected void subscribeActual(Observer observer) { + observers.add(observer); + observer.onSubscribe(Disposables.fromRunnable(() -> System.out.println("Disposed"))); + } +} diff --git a/rxjava/src/test/java/com/baeldung/rxjava/RxRelayTest.java b/rxjava/src/test/java/com/baeldung/rxjava/RxRelayTest.java new file mode 100644 index 0000000000..091e4df138 --- /dev/null +++ b/rxjava/src/test/java/com/baeldung/rxjava/RxRelayTest.java @@ -0,0 +1,120 @@ +package com.baeldung.rxjava; + +import com.jakewharton.rxrelay2.BehaviorRelay; +import com.jakewharton.rxrelay2.PublishRelay; +import com.jakewharton.rxrelay2.ReplayRelay; +import io.reactivex.internal.schedulers.SingleScheduler; +import io.reactivex.observers.TestObserver; +import org.junit.Test; + +import java.util.concurrent.TimeUnit; + +public class RxRelayTest { + + @Test + public void whenObserverSubscribedToPublishRelay_thenItReceivesEmittedEvents () { + PublishRelay publishRelay = PublishRelay.create(); + TestObserver firstObserver = TestObserver.create(); + TestObserver secondObserver = TestObserver.create(); + publishRelay.subscribe(firstObserver); + firstObserver.assertSubscribed(); + publishRelay.accept(5); + publishRelay.accept(10); + publishRelay.subscribe(secondObserver); + secondObserver.assertSubscribed(); + publishRelay.accept(15); + //First Observer will receive all events + firstObserver.assertValues(5, 10, 15); + //Second Observer will receive only last event + secondObserver.assertValue(15); + } + + @Test + public void whenObserverSubscribedToBehaviorRelayWithoutDefaultValue_thenItIsEmpty() { + BehaviorRelay behaviorRelay = BehaviorRelay.create(); + TestObserver firstObserver = new TestObserver<>(); + behaviorRelay.subscribe(firstObserver); + firstObserver.assertEmpty(); + } + + @Test + public void whenObserverSubscribedToBehaviorRelay_thenItReceivesDefaultValue() { + BehaviorRelay behaviorRelay = BehaviorRelay.createDefault(1); + TestObserver firstObserver = new TestObserver<>(); + behaviorRelay.subscribe(firstObserver); + firstObserver.assertValue(1); + } + + @Test + public void whenObserverSubscribedToBehaviorRelay_thenItReceivesEmittedEvents () { + BehaviorRelay behaviorRelay = BehaviorRelay.create(); + TestObserver firstObserver = TestObserver.create(); + TestObserver secondObserver = TestObserver.create(); + behaviorRelay.accept(5); + behaviorRelay.subscribe(firstObserver); + behaviorRelay.accept(10); + behaviorRelay.subscribe(secondObserver); + behaviorRelay.accept(15); + firstObserver.assertValues(5, 10, 15); + secondObserver.assertValues(10, 15); + } + @Test + public void whenObserverSubscribedToReplayRelay_thenItReceivesEmittedEvents () { + ReplayRelay replayRelay = ReplayRelay.create(); + TestObserver firstObserver = TestObserver.create(); + TestObserver secondObserver = TestObserver.create(); + replayRelay.subscribe(firstObserver); + replayRelay.accept(5); + replayRelay.accept(10); + replayRelay.accept(15); + replayRelay.subscribe(secondObserver); + firstObserver.assertValues(5, 10, 15); + secondObserver.assertValues(5, 10, 15); + + } + + @Test + public void whenObserverSubscribedToReplayRelayWithLimitedSize_thenItReceivesEmittedEvents () { + ReplayRelay replayRelay = ReplayRelay.createWithSize(2); + TestObserver firstObserver = TestObserver.create(); + replayRelay.accept(5); + replayRelay.accept(10); + replayRelay.accept(15); + replayRelay.accept(20); + replayRelay.subscribe(firstObserver); + + firstObserver.assertValues(15, 20); + + } + + + @Test + public void whenObserverSubscribedToReplayRelayWithMaxAge_thenItReceivesEmittedEvents () throws InterruptedException { + ReplayRelay replayRelay = ReplayRelay.createWithTime(2000, TimeUnit.MILLISECONDS, new SingleScheduler()); + TestObserver firstObserver = TestObserver.create(); + replayRelay.accept(5); + replayRelay.accept(10); + replayRelay.accept(15); + replayRelay.accept(20); + Thread.sleep(3000); + replayRelay.subscribe(firstObserver); + firstObserver.assertEmpty(); + } + + @Test + public void whenTwoObserversSubscribedToRandomRelay_thenOnlyOneReceivesEvent () { + RandomRelay randomRelay = new RandomRelay(); + TestObserver firstObserver = TestObserver.create(); + TestObserver secondObserver = TestObserver.create(); + randomRelay.subscribe(firstObserver); + randomRelay.subscribe(secondObserver); + randomRelay.accept(5); + if(firstObserver.values().isEmpty()) { + secondObserver.assertValue(5); + } else { + firstObserver.assertValue(5); + secondObserver.assertEmpty(); + } + } +} +