【问题标题】:Akka stream stops after one elementAkka 流在一个元素后停止
【发布时间】:2016-06-29 01:15:25
【问题描述】:

我的 akka 流在单个元素之后停止。这是我的直播:

val firehoseSource = Source.actorPublisher[FirehoseActor.RawTweet](
  FirehoseActor.props(
    auth = ...
  )
)

val ref = Flow[FirehoseActor.RawTweet]
  .map(r => ResponseParser.parseTweet(r.payload))
  .map { t => println("Received: " + t); t }
  .to(Sink.onComplete({
    case Success(_) => logger.info("Stream completed")
    case Failure(x) => logger.error(s"Stream failed: ${x.getMessage}")
  }))
  .runWith(firehoseSource)

FirehoseActor 连接到 Twitter firehose 并将消息缓冲到队列中。当actor收到Request消息时,它takes下一个元素并返回它:

def receive = {
  case Request(_) =>
    logger.info("Received request for next firehose element")
    onNext(RawTweet(queue.take()))
} 

问题是只有一条推文被打印到控制台。该程序不会退出或抛出任何错误,而且我已经散布了日志语句,但没有打印出来。

我认为水槽会继续施加压力以拉动元素,但似乎并非如此,因为Sink.onComplete 中的任何消息都没有被打印出来。我也尝试使用Sink.ignore,但也只打印了一个元素。 Actor 中的日志消息也只打印一次。

我需要使用什么水槽来让它无限期地拉动元素?

【问题讨论】:

    标签: akka-stream


    【解决方案1】:

    啊,我应该尊重我的演员totalDemand。这解决了问题:

    def receive = {
      case Request(_) =>
        logger.info("Received request for next firehose element")
        while (totalDemand > 0) {
          onNext(RawTweet(queue.take()))
        }
    

    我期待为流中的每个元素收到一个Request,但显然每个流都会发送一个Request

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-07-06
      • 2020-11-19
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多