目前还不清楚您想从响应中得到什么。 stream() 返回一个 Future,因此如果你在它上面 flatMap,你还需要从处理程序返回一个未来。你到底想在那里返回什么?例如,如果您想获取带有整个响应正文的Future[String],您可以使用runReduce(_ + _):
val result: Future[String] = ws.url(url).stream()
.flatMap(response => response.body
.via(framing)
.map(_.utf8String)
.map(_ + "\n")
.runReduce(_ + _)
)
runReduce(f: (U, U) => U) 返回一个Future[U],也就是说,在你的情况下它是Future[String]。如果你想通过其他函数分别处理传入流的每个元素,你可以使用runForeach:
ws.url(url).stream()
.flatMap( response => response.body
.via(framing)
.map(_.utf8String)
.map(_ + "\n")
.runForeach(s => handleString(s))
)
如果没有关于您想要做什么的更多细节,很难提供更具体的答案。
更新:如果你想限制来自外部服务器的消息,你可以使用内置的throttle 组合器:
val result: Future[Source[String, _]] = ws.url(url)
.stream()
.map { response =>
response.body
.via(framing)
.map(_.utf8String)
.map(_ + "\n")
.throttle(10, 1.second, 10, ThrottleMode.Shaping)
}
这里,result 未来将包含Strings 的流,当具体化时,每秒最多会产生 10 个元素,必要时会反压。您可以在我上面链接的文档中找到更多信息。