How to dynamically update an RX Observable?

Viewed 1421

(Working in RxKotlin and RxJava, but using metacode for simplicity)

Many Reactive Extensions guides begin by creating an Observable from already available data. From The introduction to Reactive Programming you've been missing, it's created from a single string

var soureStream= Rx.Observable.just('https://api.github.com/users');

Similarly, from the frontpage of RxKotlin, from a populated list

val list = listOf(1,2,3,4,5)
list.toObservable()     

Now consider a simple filter that yields an outStream,

var outStream = sourceStream.filter({x > 3})

In both guides the source events are declared apriori. Which means the timeline of events has some form

source: ----1,2,3,4,5-------
out:    --------------4,5---

How can I modify sourceStream to become more of a pipeline? In other words, no input data is available during sourceStream creation? When a source event becomes available, it is immediately processed by out:

source: ---1--2--3-4---5-------
out:    ------------4---5-------

I expected to find an Observable.add() for dynamic updates

var sourceStream = Observable.empty()
var outStream = sourceStream.filter({x>3})

//print each element as its added 
sourceStream .subscribe({println(it)})
outStream.subscribe({println(it)})

for i in range(5):
    sourceStream.add(i)

Is this possible?

1 Answers
Related