【问题标题】:How to complete a stream of numbers as CSV values in Akka HTTP?如何在 Akka HTTP 中完成数字流作为 CSV 值?
【发布时间】:2021-07-20 06:29:04
【问题描述】:

我有一个 Akka 源作为数字流,实现为:

Source(Stream(1, 2, 3, 4, 5))

我正在尝试利用 Akka HTTP 中的 Akka 流式传输支持以逗号分隔值的形式返回流式响应。

我已经关注akka doc 的简单 csv 源流,提出了以下实现:

implicit val csvFormat = Marshaller.strict[Int, ByteString] { res =>
    Marshalling.WithFixedContentType(ContentTypes.`text/csv(UTF-8)`, () => {
    ByteString(List(res).mkString(","))
    })
}

implicit val streamingSupport: CsvEntityStreamingSupport = EntityStreamingSupport.csv()

complete(Source(Stream(1, 2, 3, 4, 5)))

但显然这不是 CSV 实体流支持对我的目的的正确用例。这会导致每个数字都在新行中流式传输。

但这不是我想要的。我想以逗号分隔的列表形式回复,例如 1,2,3,4,5

如何使用 Akka HTTP 中的流式支持来实现这一点?

【问题讨论】:

  • 这不是 text/csv(UTF-8) - datatracker.ietf.org/doc/html/rfc4180#section-2 内容类型的有效内容。您能否在不添加任何实现细节的情况下用非常简单的语言陈述您的需求?
  • Object akka.http.scaladsl.model.ContentType 将其定义为有效值。
  • 简单来说,我需要将数字列表作为逗号分隔值而不是换行符分隔。
  • 因为您只是在流式传输单个整数......它实际上是一个 text/plain 响应。
  • 将其更改为 text/plain 会导致 curl: (18) transfer closed with outstanding read data remaining

标签: scala akka akka-stream akka-http


【解决方案1】:

该换行符由EntityStreamingSupport.csv() 中定义的流渲染器添加。

我们需要定义我们自己的自定义EntityStreamingSupport 才能使其工作。

val route =
  path("test") {
    val responseSource: Source[Int, NotUsed] =
      Source.fromIterator(() => Stream(1, 2, 3, 4, 5).iterator)

    val byteStringSource: Source[ByteString, NotUsed] =
      responseSource.map(i => ByteString(i.toString))

    val streamingSource =
      byteStringSource.map(bs => HttpEntity(ContentTypes.`text/plain(UTF-8)`, bs))

    implicit val streamingSupport =
      EntityStreamingSupport.csv(maxLineLength = 16 * 1024)
        .withSupported(ContentTypeRange(ContentTypes.`text/plain(UTF-8)`))
        .withContentType(ContentTypes.`text/plain(UTF-8)`)
        .withFramingRenderer(Flow[ByteString].map(bs => bs ++ ByteString(",")))

    complete((streamingSource))
  }
curl localhost:8080/test -v 
*   Trying 127.0.0.1...
* TCP_NODELAY set
* Connected to localhost (127.0.0.1) port 8080 (#0)
> GET /test HTTP/1.1
> Host: localhost:8080
> User-Agent: curl/7.64.1
> Accept: */*
> 
< HTTP/1.1 200 OK
< Server: akka-http/10.2.4
< Date: Tue, 20 Jul 2021 07:50:46 GMT
< Transfer-Encoding: chunked
< Content-Type: text/plain; charset=UTF-8
< 
* Connection #0 to host localhost left intact
1,2,3,4,5,* Closing connection 0

编辑:为了消除那个逗号,我们可以使用窗口黑客。

val route =
  path("test") {
    val responseSource: Source[Int, NotUsed] =
      Source.fromIterator(() => Stream(1, 2, 3, 4, 5).iterator)

    val startByteString = ByteString("$start$")

    val byteStringSource: Source[ByteString, NotUsed] =
        responseSource.map(i => ByteString(i.toString)).prepend(Source.single(startByteString))

    val streamingSource =
      byteStringSource.map(bs => HttpEntity(ContentTypes.`text/plain(UTF-8)`, bs))

    implicit val streamingSupport =
      EntityStreamingSupport.csv(maxLineLength = 16 * 1024)
        .withSupported(ContentTypeRange(ContentTypes.`text/plain(UTF-8)`))
        .withContentType(ContentTypes.`text/plain(UTF-8)`)
        .withFramingRenderer(
          Flow[ByteString].sliding(2, 1)
            .map { bsSeq =>
              if (startByteString.equals(bsSeq(0))) {
                // first int; no need for comma
                bsSeq(1)
              } else {
                // not first int; add comma
                ByteString(",") ++ bsSeq(1)
              }

            }
        )

    complete((streamingSource))
  }
curl localhost:8080/test -v 
*   Trying 127.0.0.1...
* TCP_NODELAY set
* Connected to localhost (127.0.0.1) port 8080 (#0)
> GET /test HTTP/1.1
> Host: localhost:8080
> User-Agent: curl/7.64.1
> Accept: */*
> 
< HTTP/1.1 200 OK
< Server: akka-http/10.2.4
< Date: Tue, 20 Jul 2021 08:28:05 GMT
< Transfer-Encoding: chunked
< Content-Type: text/plain; charset=UTF-8
< 
* Connection #0 to host localhost left intact
1,2,3,4,5* Closing connection 0

【讨论】:

  • 谢谢@sarveshseri。这似乎工作正常!自从您指定maxLineLength = 16 * 1024以来,我唯一担心的是这是否适用于无限数字流?
  • 我认为 http 并不是真正适合“无限”流媒体的媒体。有很多设置会在某个时间点终止请求。您应该为此使用网络套接字。
  • 似乎它与消息 curl: (18) transfer closed with outstanding read data remaining 中断,然后继续进行其余的流式传输。
  • 如果不要求提供以逗号分隔的数字流。我在问题中提出的解决方案应该适用于无限的数字流,对吧?同样适用于 SSE,对吗?
  • 没有。在所有服务器实现中,有很多“限制器”(如超时、最大内容长度等)应用于 http 请求,您可以调整限制器但不能删除它们。并且在生产环境中使用“近乎无限”的 http 请求是不可取的。
【解决方案2】:

你错过了一个更高级别的课程。它不能是一个 Int 并组成一个多列的行。创建一个如下MyBO 的类就可以了。

    case class MyBO(a: Int, b: Int, c: Int, d: Int, e: Int)
    
    implicit val myBOAsCsv = Marshaller.strict[MyBO, ByteString] { t =>
      Marshalling.WithFixedContentType(ContentTypes.`text/csv(UTF-8)`, () => {
        ByteString(List(t.a, t.b, t.c, t.d, t.e).mkString(","))
        })
      }

    implicit val csvStreaming = EntityStreamingSupport.csv()

    val route: Route =
      path("ping") {
        get {
          complete(Source(Stream(MyBO(1,3,5,7,11))))
        }
      }
 curl localhost:8080/ping -v                                                                                                                               [3f8df8e]
*   Trying 127.0.0.1...
* TCP_NODELAY set
* Connected to localhost (127.0.0.1) port 8080 (#0)
> GET /ping HTTP/1.1
> Host: localhost:8080
> User-Agent: curl/7.64.1
> Accept: */*
>
< HTTP/1.1 200 OK
< Server: akka-http/10.2.4
< Date: Tue, 20 Jul 2021 07:10:18 GMT
< Transfer-Encoding: chunked
< Content-Type: text/csv; charset=UTF-8
<
1,3,5,7,11
* Connection #0 to host localhost left intact
* Closing connection 0

【讨论】:

  • 感谢您的回答!但我认为我没有很清楚地表达我的问题陈述。对此表示歉意。这个想法是能够流式传输无限的数字流。我不认为使用案例类会有所帮助。
  • 没问题:)。很高兴看到你的答案
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-07-08
  • 2017-01-30
  • 1970-01-01
  • 2018-02-16
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多