【问题标题】:how does onNext is called in scala rx在scala rx中如何调用onNext
【发布时间】:2016-08-08 08:44:48
【问题描述】:

我试图了解 observable 的工作原理。这是我的代码。

def make: Stream[Int] = {
    Stream.cons(scala.util.Random.nextInt(10), {
      println("Making ..")
      Thread.sleep(1000)
      make
    })
  }

  val y = Observable.from(make)

  y.foreach(a => println(a))

emit 将每 1 秒产生一次新值。我正在制作一个可观察的。 for 每个循环将永远打印新生成的值。

据我了解,a=>println(a) 是一个回调值,在 rx observable 中称为 onNext(t)。

我想弄清楚的是它是如何粘在生产者身上的,所以当产生新值时,在哪里调用 onNext。我研究了一段时间的 rx 代码,但仍然无法弄清楚。

谢谢。

【问题讨论】:

    标签: scala system.reactive reactive-programming


    【解决方案1】:

    似乎订阅了流调用 rx.Observable 的

    Subscription subscribe(Subscriber<? super T> subscriber) 
    

    这将调用 IterableProducer.request 方法 其中会有这段代码。

        if (n == Long.MAX_VALUE && REQUESTED_UPDATER.compareAndSet(this, 0, Long.MAX_VALUE)) {
            // fast-path without backpressure
            while (it.hasNext()) {
                if (o.isUnsubscribed()) {
                    return;
                }
                o.onNext(**it.next()**);
            }
            if (!o.isUnsubscribed()) {
                o.onCompleted();
            }
        }
    

    onNext 是观察者代码,它的输入将是生产者 it.next。就是这样粘的。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多