【问题标题】:Transforming Slick Streaming data and sending Chunked Response using Akka Http使用 Akka Http 转换 Slick Streaming 数据并发送分块响应
【发布时间】:2018-06-08 01:55:24
【问题描述】:

目的是从数据库中流式传输数据,对这块数据执行一些计算(此计算返回某个案例类的 Future)并将这些数据作为分块响应发送给用户。目前,我能够在不执行任何计算的情况下流式传输数据并发送响应。但是,我无法执行此计算然后流式传输结果。

这是我实现的路线。

def streamingDB1 =
path("streaming-db1") {
  get {
    val src = Source.fromPublisher(db.stream(getRds))
    complete(src)
  }
}

函数 getRds 返回映射到案例类的表的行(使用 slick)。现在考虑将每一行作为输入并返回另一个案例类的 Future 的函数 compute。像

def compute(x: Tweet) : Future[TweetNew] = ?

如何在变量 src 上实现这个函数,并将这个计算的分块响应(作为流)发送给用户。

【问题讨论】:

    标签: scala akka slick akka-stream akka-http


    【解决方案1】:

    您可以使用mapAsync 转换源代码:

    val src =
      Source.fromPublisher(db.stream(getRds))
            .mapAsync(parallelism = 3)(compute)
    
    complete(src)
    

    根据需要调整并行度。


    请注意,您可能需要配置Slick documentation 中提到的一些设置:

    注意:某些数据库系统可能需要以某种方式设置会话参数以支持流式传输,而无需在客户端的内存中一次缓存所有数据。例如,PostgreSQL 需要.withStatementParameters(rsType = ResultSetType.ForwardOnly, rsConcurrency = ResultSetConcurrency.ReadOnly, fetchSize = n)(具有所需的页面大小n)和.transactionally 才能进行正确的流式传输。

    因此,例如,如果您使用的是 PostgreSQL,那么您的 Source 可能如下所示:

    val src =
      Source.fromPublisher(
        db.stream(
          getRds.withStatementParameters(
            rsType = ResultSetType.ForwardOnly,
            rsConcurrency = ResultSetConcurrency.ReadOnly,
            fetchSize = 10
          ).transactionally
        )
      ).mapAsync(parallelism = 3)(compute)
    

    【讨论】:

    • 这不起作用。我运行 curl 命令来命中端点。但是连接会关闭。
    • @user3294786 听起来好像在关闭之前没有正确等待数据
    • @StanislavPalatnik 可能在这里有所作为;我建议添加一个 .log() 阶段以查看元​​素实际上是按预期发送的(而不仅仅是完成),确保也设置 akka.loglevel = DEBUG
    • @user3294786:您将此标记为答案。那么它现在对你有用吗?
    • @StanislavPalatnik 是的,它奏效了。我在不同的数据集上尝试了这个,它工作得很好。需要调试它之前失败的原因。
    【解决方案2】:

    您需要有一种方法来编组 TweetNew,并且如果您发送长度为 0 的块,客户端可能会关闭连接。

    此代码适用于 curl:

    case class TweetNew(str: String)
    
    def compute(string: String) : Future[TweetNew] = Future {
      TweetNew(string)
    }
    
    val route = path("hello") {
      get {
        val byteString: Source[ByteString, NotUsed] = Source.apply(List("t1", "t2", "t3"))
          .mapAsync(2)(compute)
          .map(tweet => ByteString(tweet.str + "\n"))
        complete(HttpEntity(ContentTypes.`text/plain(UTF-8)`, byteString))
      }
    }
    

    【讨论】:

    • "如果你发送一个长度为 0 的块,客户端可能会关闭连接。" -- 我相信 Akka 足够聪明,可以为您从流中过滤掉这些,因为它们不能表示为 HTTP 块。
    猜你喜欢
    • 2016-01-12
    • 2016-01-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-04-01
    • 2011-09-11
    • 2022-11-27
    • 1970-01-01
    相关资源
    最近更新 更多