Monix: InputStreamObservable does not support multiple subscribers

I am trying to split Observable (String, Date) into two different Observables and pin them together as follows

import monix.execution.Scheduler.Implicits.global
val x = Observable.fromIterator((0 to 10).map(i => (s"a $i", s"b $i")).toIterator)

val y = Observable.toReactive(x)

val fileStream = Observable.fromReactivePublisher(y).mapAsync(5)(a => Task{println(a._1); a._1})
val dateStream = Observable.fromReactivePublisher(y).mapAsync(5)(a => Task{println(a._2); a._2})

fileStream.zip(dateStream)
  .map(println)
  .subscribe()

But I get the following exception

monix.reactive.exceptions.MultipleSubscribersException: InputStreamObservable does not support multiple subscribers
    at monix.reactive.exceptions.MultipleSubscribersException$.build(MultipleSubscribersException.scala:51)
    at monix.reactive.internal.builders.IteratorAsObservable.unsafeSubscribeFn(IteratorAsObservable.scala:42)
    at monix.reactive.Observable$$anon$6.subscribe(Observable.scala:155)
    at monix.reactive.internal.builders.ReactiveObservable.unsafeSubscribeFn(ReactiveObservable.scala:38)
    at monix.reactive.internal.operators.MapAsyncParallelObservable.unsafeSubscribeFn(MapAsyncParallelObservable.scala:60)
    at monix.reactive.internal.builders.Zip2Observable.unsafeSubscribeFn(Zip2Observable.scala:158)
    at monix.reactive.Observable$$anon$5.unsafeSubscribeFn(Observable.scala:139)
    at monix.reactive.Observable$class.subscribe(Observable.scala:71)
    at monix.reactive.Observable$$anon$5.subscribe(Observable.scala:136)
    at monix.reactive.Observable$class.subscribe(Observable.scala:90)
    at monix.reactive.Observable$$anon$5.subscribe(Observable.scala:136)
    at monix.reactive.Observable$class.subscribe(Observable.scala:120)
    at monix.reactive.Observable$$anon$5.subscribe(Observable.scala:136)
    at monix.reactive.Observable$class.subscribe(Observable.scala:112)
    at monix.reactive.Observable$$anon$5.subscribe(Observable.scala:136)
+4
source share
2 answers

Transformation to / from reactive is mandatory?

One way to fix this is val x = Observable.fromIterable((0 to 10).map(i => (s"a $i", s"b $i"))), but it will throw an OutOfMemoryError for infinity streams.

Another way is to use .multicast(Pipe.publish[])and then obs.connect()down the code:

import monix.execution.Scheduler.Implicits.global
val x = Observable.fromIterator((0 to 10).map(i => (s"a $i", s"b $i")).iterator)

val y = Observable.toReactive(x)
val obsY = Observable.fromReactivePublisher(y)
val connectY = obsY.multicast(Pipe.publish[(String, String)])

val fileStream = connectY.mapAsync(5)(a => Task{println(a._1); a._1})
val dateStream = connectY.mapAsync(5)(a => Task{println(a._2); a._2})

fileStream.zip(dateStream)
  .map(println)
  .subscribe()

connectY.connect()

Thread.sleep(5000)
+2
source

sergei-shubin Observable "" , monix.reactive.Observable[R]):monix.reactive.Observable[R] rel="nofollow noreferrer"> publishSelector multicast. :

val x = Observable.fromIterator((0 to 10).map(i => (s"a $i", s"b $i")).toIterator)

val zipped = x.publishSelector { o =>
  val fileStream = o.mapParallelUnordered(5)(a => Task{println(a._1); a._1})
  val dateStream = o.mapParallelUnordered(5)(a => Task{println(a._2); a._2})

  fileStream.zip(dateStream)
}

zipped
  .map(println)
  .subscribe()
0

Source: https://habr.com/ru/post/1693297/


All Articles