【发布时间】: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