【问题标题】:How to convert reactive Publisher to simple Stream in Scala?如何在 Scala 中将反应式发布者转换为简单的流?
【发布时间】:2019-03-01 14:00:08
【问题描述】:

是否可以将org.reactivestreams.Publisher 实例转换为scala.Stream?如果可以,怎么做?

【问题讨论】:

    标签: scala reactive-streams


    【解决方案1】:

    以下内容对您有用吗?

    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 我对“反应性流”一无所知。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-01-31
    • 2020-05-13
    • 2011-03-13
    相关资源
    最近更新 更多