gpt4 book ai didi

scala - 拆分 Monix Observable

转载 作者:行者123 更新时间:2023-12-01 09:32:49 32 4
gpt4 key购买 nike

我想为 monix.reactive.Observable 编写一个拆分函数.它应该拆分一个源 Observable[A]变成一对新的 (Observable[A], Observable[A]) ,基于谓词的值,针对源中的每个元素进行评估。我希望拆分独立于源 Observable 是热的还是冷的。在源是冷的情况下,新的 Observable 对也应该是冷的,而在源是热的情况下,新的 Observable 对将是热的。我想知道这样的实现是否可行,如果可以,如何实现(我在下面粘贴了一个失败的测试用例)。

签名,作为隐式类上的方法,看起来像或类似于

    /**
* Split an observable by a predicate, placing values for which the predicate returns true
* to the right (and values for which the predicate returns false to the left).
* This is consistent with the convention adopted by Either.cond.
*/
def split(p: T => Boolean)(implicit scheduler: Scheduler, taskLike: TaskLike[Future]): (Observable[T], Observable[T]) = {
splitEither[T, T](elem => Either.cond(p(elem), elem, elem))
}

目前,我有一个简单的实现,它使用源元素并将它们推送到 PublishSubject .因此,这对新的 Observable 很热门。我对冷 Observable 的测试失败了。

import monix.eval.TaskLike
import monix.execution.{Ack, Scheduler}
import monix.reactive.{Observable, Observer}
import monix.reactive.subjects.PublishSubject

import scala.concurrent.Future

object ObservableOps {

implicit class ObservableExtensions[T](o: Observable[T]) {

/**
* Split an observable by a predicate, placing values for which the predicate returns true
* to the right (and values for which the predicate returns false to the left).
* This is consistent with the convention adopted by Either.cond.
*/
def split(p: T => Boolean)(implicit scheduler: Scheduler, taskLike: TaskLike[Future]): (Observable[T], Observable[T]) = {
splitEither[T, T](elem => Either.cond(p(elem), elem, elem))
}

/**
* Split an observable into a pair of Observables, one left, one right, according
* to a determinant function.
*/
def splitEither[U, V](f: T => Either[U, V])(implicit scheduler: Scheduler, taskLike: TaskLike[Future]): (Observable[U], Observable[V]) = {
val l = PublishSubject[U]()
val r = PublishSubject[V]()

o.subscribe(new Observer[T] {
override def onNext(elem: T): Future[Ack] = {
f(elem) match {
case Left(u) => l.onNext(u)
case Right(v) => r.onNext(v)
}
}

override def onError(ex: Throwable): Unit = {
l.onError(ex)
r.onError(ex)
}

override def onComplete(): Unit = {
l.onComplete()
r.onComplete()
}
})

(l, r)
}
}
}

//////////


import ObservableOps._

import monix.execution.Scheduler.Implicits.global
import monix.reactive.Observable
import monix.reactive.subjects.PublishSubject

import org.scalatest.FlatSpec
import org.scalatest.Matchers._
import org.scalatest.concurrent.ScalaFutures._

class ObservableOpsSpec extends FlatSpec {

val isEven: Int => Boolean = _ % 2 == 0

"Observable Ops" should "split a cold observable" in {
val o = Observable(1, 2, 3, 4, 5)

val (l, r) = o.split(isEven)

l.toListL.runToFuture.futureValue shouldBe List(1, 3, 5)
r.toListL.runToFuture.futureValue shouldBe List(2, 4)
}

"Observable Ops" should "split a hot observable" in {
val o = PublishSubject[Int]()

val (l, r) = o.split(isEven)
val lbuf = l.toListL.runToFuture
val rbuf = r.toListL.runToFuture

Observable.fromIterable(1 to 5).mapEvalF(i => o.onNext(i)).subscribe()
o.onComplete()

lbuf.futureValue shouldBe List(1, 3, 5)
rbuf.futureValue shouldBe List(2, 4)
}
}

我希望上面的两个测试用例都能通过,但 "Observable Ops" should "split a cold observable"正在失败。

编辑:工作代码

通过两个测试用例的实现如下:

import monix.execution.Scheduler
import monix.reactive.Observable

object ObservableOps {

implicit class ObservableExtension[T](o: Observable[T]) {

/**
* Split an observable by a predicate, placing values for which the predicate returns true
* to the right (and values for which the predicate returns false to the left).
* This is consistent with the convention adopted by Either.cond.
*/
def split(
p: T => Boolean
)(implicit scheduler: Scheduler): (Observable[T], Observable[T]) = {
splitEither[T, T](elem => Either.cond(p(elem), elem, elem))
}

/**
* Split an observable into a pair of Observables, one left, one right, according
* to a determinant function.
*/
def splitEither[U, V](
f: T => Either[U, V]
)(implicit scheduler: Scheduler): (Observable[U], Observable[V]) = {

val oo = o.map(f)

val l = oo.collect {
case Left(u) => u
}

val r = oo.collect {
case Right(v) => v
}

(l, r)
}
}
}

最佳答案

class ObservableOpsSpec extends FlatSpec {

val isEven: Int => Boolean = _ % 2 == 0

"Observable Ops" should "split a cold observable" in {
val o = Observable(1, 2, 3, 4, 5)
val o2 = o.publish
val (l, r) = o2.split(isEven)

val x= l.toListL.runToFuture
val y = r.toListL.runToFuture
o2.connect()
x.futureValue shouldBe List(1, 3, 5)
y.futureValue shouldBe List(2, 4)
}

"Observable Ops" should "split a hot observable" in {
val o = PublishSubject[Int]()

val (l, r) = o.split(isEven)
val lbuf = l.toListL.runToFuture
val rbuf = r.toListL.runToFuture

Observable.fromIterable(1 to 5).mapEvalF(i => o.onNext(i)).subscribe()
o.onComplete()

lbuf.futureValue shouldBe List(1, 3, 5)
rbuf.futureValue shouldBe List(2, 4)
}
}

关于scala - 拆分 Monix Observable,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/57909519/

32 4 0
Copyright 2021 - 2024 cfsdn All Rights Reserved 蜀ICP备2022000587号
广告合作:1813099741@qq.com 6ren.com