【问题标题】:Rx (RxKotlin) - rightGroupJoin using groupJoin - merge / combine two observables of different typesRx (RxKotlin) - rightGroupJoin 使用 groupJoin - 合并/组合两个不同类型的 observable
【发布时间】:2019-04-17 12:18:07
【问题描述】:

在苦苦挣扎了几天之后,为了一件看似简单的事情,我来找你们了:)

想法很简单。我有两个流/可观察对象,“左”和“右”。 我希望“右”中的项目缓冲/收集/聚合到“左”中的“当前”项目。
因此,“左”中的每个项目都定义了一个新的“窗口”,而所有“右”项目都将绑定到该窗口,直到发出新的“左”项目。所以,可视化:

任务:
'左'     : |- A - - - - - B - - C - - - -|
'正确' : |- 1 - 2 - 3 -4 - 5 - 6 - - -|
'结果' : |- - - - - - - -x - - -y - - - -z| (Pair<Left, List<Right>>)
在哪里:A,1B,4 (so x) ; C(所以 y)同时发出
所以:     x = Pair(A, [1,2,3]),  y = Pair(B , [4, 5])
并且:'right' & 'result' 完成/当 'left' 完成时终止
所以:    z = Pair( C, [6]) - 作为 'left' 完成的结果发出

----
编辑 2 - 最终解决方案!
为了将“右”项目与下一个“左”而不是前一个聚合,我将代码更改为更短/更简单的代码:

fun <L, R> Observable<L>.rightGroupJoin(right: Observable<R>): Observable<Pair<L, List<R>>> {
    return this.share().run {
        zipWith(right.buffer(this), BiFunction { left, rightList ->
            Pair(left, rightList)
        })
    }
}  

编辑 1 - 初始解决方案!
摘自下面@Mark(已接受)的答案,这就是我想出的。
它被分成更小的方法,因为我还使用multiRightGroupJoin() 来加入尽可能多的(正确的)流。

fun <T, R> Observable<T>.rightGroupJoin(right: Observable<R>): Observable<Pair<T, List<R>>> {
    return this.share().let { thisObservable ->    //use 'share' to avoid multi-subscription complications, e.g. multi calls to **preceding** doOnComplete
        thisObservable.flatMapSingle { t ->        //treat each 'left' as a Single
            bufferRightOnSingleLeft(thisObservable, t, right)
        }
    }
}

地点:

private fun <T, R> bufferRightOnSingleLeft(left: Observable<*>, leftSingleItem: T, right: Observable<R>)
    : Single<Pair<T, MutableList<R>>> {

    return right.buffer(left)                              //buffer 'right' until 'left' onNext() (for each 'left' Single) 
        .map { Pair(leftSingleItem, it) }
        .first(Pair(leftSingleItem, emptyList()))   //should be only 1 (list). THINK firstOrError
}  

----

到目前为止我得到了什么
经过大量阅读并了解到不知何故没有开箱即用的实现,我决定使用groupJoin,主要使用this link,如下所示:(这里有很多问题和需要改进的地方,不要'不要使用这个代码)

private fun <T, R> Observable<T>.rightGroupJoin(right: Observable<R>): Observable<Pair<T, List<R>>> {

var thisCompleted = false //THINK is it possible to make the groupJoin complete on the left(this)'s onComplete automatically?
val thisObservable = this.doOnComplete { thisCompleted = true }
        .share() //avoid weird side-effects of multiple onSubscribe calls

//join/attach 'right/other' stream to windows (buffers), starting and ending on each 'this/left' onNext
return thisObservable.groupJoin(

    //bind 'right/other' stream to 'this/left'
    right.takeUntil { thisCompleted }//have an onComplete rule THINK add share() at the end?

    //define when windows start/end ('this/left' onNext opens new window and closes prev)
    , Function<T, ObservableSource<T>> { thisObservable }

    //define 'right/other' stream to have no windows/intervals/aggregations by itself
    // -> immediately bind each emitted item to a 'current' window(T) above
    , Function<R, ObservableSource<R>> { Observable.empty() }

    //collect the whole 'right' stream in 'current' ('left') window
    , BiFunction<T, Observable<R>, Single<Pair<T, List<R>>>> { t, rObs ->
        rObs.collect({ mutableListOf<R>() }) { acc, value ->
            acc.add(value)
        }.map { Pair(t, it.toList()) }

    }).mergeAllSingles()
}  

我也使用类似的用法来创建timedBuffer() - 与buffer(timeout) 相同,但在每个缓冲区上都有一个时间戳(List)以了解它何时开始。基本上通过在Observable.interval(timeout)(作为“左”)上运行相同的代码

问题/问题(从最简单到最难)

  1. 这是做类似事情的最佳方式吗?这不是矫枉过正吗?
  2. 当“左”完成时,是否有更好的方法(必须)来完成“结果”(和“右”)?没有这个丑陋的布尔逻辑?
  3. 这种用法似乎弄乱了 rx 的顺序。请参阅下面的代码并打印:

    leftObservable
    .doOnComplete {
        log("doOnComplete - before join")
     }
    .doOnComplete {
        log("doOnComplete 2 - before join")
     }
    .rightGroupJoin(rightObservable)
    .doOnComplete {
        log("doOnComplete - after join")
     }
    

打印(有时!看起来像竞争条件)以下内容:
doOnComplete - before join
doOnComplete - after join
doOnComplete 2 - before join

  1. 上述代码第一次运行时,doOnComplete - after join 没有被调用,第二次被调用两次。第三次就像第一次,第四次就像第二次,等等...
    3,4 都使用此代码运行。可能与 subscribe {} 用法有关?请注意,我不持有一次性用品。 这个流结束了,因为我 GC 'left' observable

    leftObservable.subscribeOn().observeOn()
    .doOnComplete{log...}
    .rightGroupJoin()
    .doOnComplete{log...}
    .subscribe {}  
    

注意1:在mergeAllSingles() 之后添加.takeUntil { thisCompleted } 似乎可以修复#4。

注意2:在使用此方法加入多个流并应用'Note1'之后,很明显onComplete(在groupJoin()调用之前!!!)将被调用的次数与'right' Observables一样多,可能这意味着原因是right.takeUntil { thisCompleted },关闭“正确”流真的很重要吗?

Note3:关于 Note1,它似乎与 takeUntil 与 takeWhile 非常相关。使用 takeWhile 降低了 doOnComplete 调用,这在某种程度上是合乎逻辑的。仍在努力解决问题。

  1. 除了在 groupJoin * rightObservablesCount 上运行 zip 之外,您能想到一个 multiGroupJoin,或者在我们的例子中是 multiRightGroupJoin?

请随便问。我知道事实上我对订阅/一次性使用和手册 onComplete 的使用不是这样,我只是不确定是什么..

【问题讨论】:

    标签: rxjs rx-java reactive-programming rx-kotlin rx-kotlin2


    【解决方案1】:

    这样简单的事情应该可以工作:

    @JvmStatic
    fun main(string: Array<String>) {
        val left = PublishSubject.create<String>()
        val right = PublishSubject.create<Int>()
    
        left.flatMapSingle { s ->  right.buffer(left).map { Pair(s, it) }.firstOrError() }
                .subscribe{ println("Group : Letter : ${it.first}, Elements : ${it.second}") }
    
    
        left.onNext("A")
        right.onNext(1)
        right.onNext(2)
        right.onNext(3)
        left.onNext("B")
        right.onNext(4)
        right.onNext(5)
        left.onNext("C")
        right.onNext(6)
        left.onComplete()
    }
    

    输出:

    Group : Letter : A, Elements : [1, 2, 3]
    Group : Letter : B, Elements : [4, 5]
    Group : Letter : C, Elements : [6]
    

    您感兴趣的Observable 是左边,所以订阅它。然后只需通过左 observable 的下一个发射或完成来缓冲右。您只对每个上游左发射的单个结果感兴趣,所以只需使用flatMapSingle。我选择了firstOrError(),但显然可以有一个默认项或其他错误处理,甚至是flatMapMaybe 加上firstElement()

    编辑

    OP 进行了进一步的问答,并发现原始问题和上述解决方案用前一个左发射缓冲右值,直到下一个左发射(如上),不是必需的行为。新要求的行为是将右值缓冲到 NEXT 左发射,如下所示:

    @JvmStatic
        fun main(string: Array<String>) {
            val left = PublishSubject.create<String>()
            val right = PublishSubject.create<Int>()
    
    
            left.zipWith (right.buffer(left), 
                    BiFunction<String, List<Int>, Pair<String, List<Int>>> { t1, t2 -> Pair(t1, t2)
            }).subscribe { println("Group : Letter : ${it.first}, Elements : ${it.second}") }
    
            left.onNext("A")
            right.onNext(1)
            right.onNext(2)
            right.onNext(3)
            left.onNext("B")
            right.onNext(4)
            right.onNext(5)
            left.onNext("C")
            right.onNext(6)
            left.onComplete()
        }
    

    这会产生不同的最终结果,因为左值与之前的右值压缩在一起,直到下一个左发射(反向)。

    输出:

    Group : Letter : A, Elements : []
    Group : Letter : B, Elements : [1, 2, 3]
    Group : Letter : C, Elements : [4, 5]
    

    【讨论】:

    • 看起来简单而完美。很快就会实施,看看是否合适,谢谢!
    • val left = PublishSubject.create&lt;String&gt;().publish().refCount 有意义吗?否则任何前面的 doOnComplete 都会被调用两次(我们在你的代码中传递了 'left' 两次)
    • 我使用的主题只是用于推送事件,您可以轻松拥有发射器或自定义可观察对象。如果使用仅在订阅时发出的ConnectableObservable 有意义,那么在“热”可观察对象上就可以了,因为这是 IMO 的实现细节问题。 doOn 运算符应该只用于副作用或日志记录,这就是我假设您使用它们的方式。
    • 是的,我用它们来记录,以进行健全性检查。多次调用它们只是觉得 rx 是一种反模式,但我真的不知道(因此提出了问题)。谢谢!
    • 经过进一步的 QA,很明显,这个解决方案在 已发出的“左”单曲上聚合了“右”,这意味着“右”项与前一个配对'left',而不是下一个。解决方案是将flatMapSingle 更改为zipWith(right.buffer(left), ..),这样当发出'left' 时,已经发出的'right' 项与之配对。您能否编辑您的答案以添加此解决方案?我仍然希望它是被接受的。还编辑了我的问题(使用新解决方案)
    【解决方案2】:

    乍一看,我会在这里使用 2 scans。示例:

    data class Result(val left: Left?, val rightList: List<Right>) {
        companion object {
            val defaultInstance: Result = Result(null, listOf())
        }
    }
    
    leftObservable.switchMap { left -> 
        rightObservable.scan(listOf()) {list, newRight -> list.plus(newRight)}
            .map { rightsList -> Result(left, rightList) }
    }
    .scan(Pair(Result.defaultInstance, Result.defaultInstance)) { oldPair, newResult -> 
        Pair(oldPair.second, newResult)
    }
    .filter { it.first != it.second }
    .map { it.first }
    

    这里唯一的问题是处理onComplete,不知道怎么做

    【讨论】:

    • 谢谢!很快就会尝试一下:)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-05-20
    • 1970-01-01
    • 2018-04-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多