【问题标题】:How to read response as Observable[String] with sttp如何使用 sttp 将响应读取为 Observable[String]
【发布时间】:2019-06-28 08:06:06
【问题描述】:

我正在使用 sttp 客户端。我想将响应解释为除以行的字符串,例如Observable[String]

这里是sttp流api:

import java.nio.ByteBuffer

import com.softwaremill.sttp._
import com.softwaremill.sttp.okhttp.monix.OkHttpMonixBackend
import monix.eval.Task
import monix.reactive.Observable

implicit val sttpBackend = OkHttpMonixBackend()

val res: Task[Response[Observable[ByteBuffer]]] = sttp
  .post(uri"someUri")
  .response(asStream[Observable[ByteBuffer]])
  .send()

那么我怎样才能得到Observable[String]

这里有一些想法:

1. 有没有一种简单的方法可以通过线观察split
2. 或者我可以从响应中得到原始的InputStream,所以我可以轻松拆分它,但我找不到使用 asStream[InputStream]
3. 之类的方法,或者只是使用 http 后端而不使用 sttp 层?

【问题讨论】:

标签: scala reactive-programming monix sttp


【解决方案1】:

您的基本问题是如何将Observable[ByteBuffer] 转换为Observable[String],其中每个String 是一行,对吗?

您可以使用方法bufferWithSelector(selector: Observable[S]): Observable[Seq[A]]。 此方法将缓冲 Observable,直到选择器 Observable 发出一个元素。

我用Ints做了一个小例子:

import monix.reactive.Observable
import monix.execution.Scheduler.Implicits.global
import scala.concurrent.duration._

val source = Observable.range(0, 1000, 1)
  .delayOnNext(100.milliseconds)

val selector = source.filter(_ % 10 == 0)

val buffered = source.bufferWithSelector(selector)
  .map(_.foldLeft("")((s, i) => s + i.toString)) // This folds the Seq[Int] into a String for display purposes

buffered.foreach(println)

Try it out!


当然,这有一个主要缺点:底层的 Observable source 将被评估两次。你可以通过修改上面的例子看到这一点:

// Start writing your ScalaFiddle code here

import monix.reactive.Observable
import monix.execution.Scheduler.Implicits.global
import scala.concurrent.duration._

val source = Observable.range(0, 1000, 1)
  .delayOnNext(100.milliseconds)
  .map {x => println(x); x}  // <------------------

val selector = source.filter(_ % 10 == 0)

val buffered = source.bufferWithSelector(selector)
  .map(_.foldLeft("")((s, i) => s + i.toString))

buffered.foreach(println)

这会将每个数字打印两次。


要解决此问题,您必须将 source Observable 转换为热 Observable:

import monix.reactive.Observable
import monix.execution.Scheduler.Implicits.global
import scala.concurrent.duration._

val source = Observable.range(0, 1000, 1)
  .delayOnNext(100.milliseconds)
  .map {x => println(x); x}
  .publish // <-----------------------------

// source is now a ConnectableObservable and will start emitting elements
// once you call source.connect()

val selector = source.filter(_ % 10 == 0)

val buffered = source.bufferWithSelector(selector)
  .map(_.foldLeft("")((s, i) => s + i.toString))

buffered.foreach(println)

source.connect() // <---------------------------

Try it out!

您唯一需要做的就是修改选择器以仅发出项目 当遇到换行时。

我建议先将Observable[ByteBuffer] 拆分为Observable[Byte](使用flatMap)以避免头痛。

【讨论】:

    猜你喜欢
    • 2021-11-28
    • 1970-01-01
    • 2023-03-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-07-27
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多