【发布时间】:2021-10-08 01:30:41
【问题描述】:
我们如何在 scala play 框架中使用 SSE?我能找到的大部分资源都是用来制作 SSE 源的。我想可靠地监听来自其他服务的 SSE 事件(使用自动连接)。最相关的文章是https://doc.akka.io/docs/alpakka/current/sse.html。我实现了这个,但这似乎不起作用(下面的代码)。还有我喜欢的事件
@Singleton
class SseConsumer @Inject()((implicit ec: ExecutionContext) {
implicit val system = ActorSystem()
val send: HttpRequest => Future[HttpResponse] = foo
def foo(x:HttpRequest) = {
try {
println("foo")
val authHeader = Authorization(BasicHttpCredentials("user", "pass"))
val newHeaders = x.withHeaders(authHeader)
Http().singleRequest(newHeaders)
}catch {
case e:Exception => {
println("Exception", e.printStackTrace())
throw e
}
}
}
val eventSource: Source[ServerSentEvent, NotUsed] =
EventSource(
uri = Uri("https://abc/v1/events"),
send,
initialLastEventId = Some("2"),
retryDelay = 1.second
)
def orderStatusEventStable() = {
val events: Future[immutable.Seq[ServerSentEvent]] =
eventSource
.throttle(elements = 1, per = 500.milliseconds, maximumBurst = 1, ThrottleMode.Shaping)
.take(10)
.runWith(Sink.seq)
events.map(_.foreach( x => {
println("456")
println(x.data)
}))
}
Future {
blocking{
while(true){
try{
Thread.sleep(2000)
orderStatusEventStable()
} catch {
case e:Exception => {
println("Exception", e.printStackTrace())
}
}
}
}
}
}
这不会给出任何异常,并且永远不会打印 println("456")。
编辑:
Future {
blocking {
while(true){
try{
Await.result(orderStatusEventStable() recover {
case e: Exception => {
println("exception", e)
throw e
}
}, Duration.Inf)
} catch {
case e:Exception => {
println("Exception", e.printStackTrace())
}
}
}
}
}
添加了等待并开始工作。一次可以阅读10条消息。但现在我面临另一个问题。 我有一个生产者,它的生产速度有时比我消耗的快,并且使用此代码我有 2 个问题:
- 我必须等到 10 条消息可用。我们怎样才能达到最大值。 10 分钟。共 0 条消息?
- 当生产率 > 消费率时,我错过了一些事件。我猜这是由于节流。我们如何使用背压来处理它?
【问题讨论】:
-
您可能需要在
Future上添加recover,因为异常会以Future.failed的形式发生,但不会出现在catch子句中
标签: scala playframework akka akka-stream server-sent-events