【问题标题】:Consuming Server Sent Events(SSE) in scala play framework with automatic reconnect在具有自动重新连接的scala play框架中使用服务器发送事件(SSE)
【发布时间】: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 个问题:

  1. 我必须等到 10 条消息可用。我们怎样才能达到最大值。 10 分钟。共 0 条消息?
  2. 当生产率 > 消费率时,我错过了一些事件。我猜这是由于节流。我们如何使用背压来处理它?

【问题讨论】:

  • 您可能需要在Future 上添加recover,因为异常会以Future.failed 的形式发生,但不会出现在catch 子句中

标签: scala playframework akka akka-stream server-sent-events


【解决方案1】:

您的代码中的问题是 events: Future 只会在流 (eventSource) 完成时完成。

我不熟悉 SSE,但在您的情况下,流可能永远不会完成,因为它总是在监听新事件。

您可以在 Akka Stream 文档中了解更多信息。

根据您想对事件执行的操作,您可以在流中 map,例如:

eventSource
  ...
  .map(/* do something */)
  .runWith(...)

基本上,您需要使用 Akka Stream Source,因为数据正在通过它,但不要等待它完成。

编辑:我没有注意到take(10),我的回答仅适用于take 不在这里的情况。您的代码应该在发送 10 个事件后工作。

【讨论】:

  • 是的,我在测试时没有等到 10 条消息,我真傻。也可以请看一下编辑并帮助我找出我面临的问题吗?
  • 请为您的每个问题打开新问题。这往往很难阅读,您将几个问题混合为一个问题。
  • 打开了一个新问题,stackoverflow.com/questions/68628414/…
猜你喜欢
  • 1970-01-01
  • 2012-07-15
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多