【发布时间】:2019-03-01 14:00:08
【问题描述】:
是否可以将org.reactivestreams.Publisher 实例转换为scala.Stream?如果可以,怎么做?
【问题讨论】:
是否可以将org.reactivestreams.Publisher 实例转换为scala.Stream?如果可以,怎么做?
【问题讨论】:
以下内容对您有用吗?
val queue: java.util.concurrent.BlockingQueue[T] = ... // TODO: choose appropriate BlockingQueue implementation
publisher.subscribe(new Subscriber[T] {
override def onNext(t: T): Unit = { queue.put(t) }
// TODO: implement other Subscriber methods
}
val stream = Stream.continually(queue.take)
【讨论】:
onNext 永远不会被调用
queue.poll 在我尝试从流中获取值时被调用。队列为空
BlockingQueue,我已经改变了答案。
take 正在阻塞,如果队列中的元素为空,它会等待该元素。队列将完成队列应该做的工作,我的意思是缓冲消息。你在上面说过onNext 永远不会被调用。我仍然不明白为什么,但在这里几乎无法提供帮助,b/c 我对“反应性流”一无所知。