【问题标题】:Convert infinite stream of finite streams to an infinite stream - Reactive X将有限流的无限流转换为无限流 - Reactive X
【发布时间】:2017-09-02 17:18:53
【问题描述】:

如何在 Reactive x 中(最好使用 RxJava 或 RxJs 中的示例)实现这一点?

a |-a-------------------a-----------a-----------a----
s1 |-x-x-x-x-x-x -| (subscribe)
s2                       |-x-x-x-x-x-| (subscribe)
s2                                               |-x-x-x-x-x-| (subscribe)
...
sn
S |-x-x-x-x-x-x-x-------x-x-x-x-x-x-x-------------x-x-x-x-x-x- (subsribe)

a 是一个无限的事件流,它触发有限流sn 的事件,每个事件都应该是无限流S 的一部分,同时能够订阅每个sn 流(以便进行求和操作),但同时保持流 S 为无限。

编辑:更具体地说,我提供了我在 Kotlin 中寻找的实现。 每 10 秒发出一个事件,该事件映射到 4 个事件的共享有限流。元流是flatMap-ed 进入正常的无限流。我利用doAfterNext 额外订阅每个有限流并打印结果。

/** Creates a finite stream with events
 * $ch-1 - $ch-4
 */
fun createFinite(ch: Char): Observable<String> =
        Observable.interval(1, TimeUnit.SECONDS)
                .take(4)
                .map({ "$ch-$it" }).share()

fun main(args: Array<String>) {

    var ch = 'A'

    Observable.interval(10, TimeUnit.SECONDS).startWith(0)
            .map { createFinite(ch++) }
            .doAfterNext {
                it
                        .count()
                        .subscribe({ c -> println("I am done. Total event count is $c") })
            }
            .flatMap { it }
            .subscribe { println("Just received [$it] from the infinite stream ") }

    // Let main thread wait forever
    CountDownLatch(1).await()
}

但是我不确定这是否是“纯 RX”方式。

【问题讨论】:

  • 这看起来像concatMap,但从问题中不清楚如何将每个事件映射到一组 N 个内部源。
  • 也许可以添加一个您迄今为止尝试过的示例,这将使我们更好地了解您要完成的工作。
  • @inf 我知道这并不理想,因为我刚刚被 rx 弄湿了。欢迎您编辑它。
  • @dev-null 真的是个笑话,我不知道rx 是什么。

标签: rxjs rx-java reactivex reactive-streams


【解决方案1】:

您没有明确说明要如何进行计数。如果你是在做总数,那么就不需要做内部订阅了:

AtomicLong counter = new AtomicLong()
Observable.interval(10, TimeUnit.SECONDS).startWith(0)
        .map { createFinite(ch++) }
        .flatMap { it }
        .doOnNext( counter.incrementAndget() )
        .subscribe { println("Just received [$it] from the infinite stream ") }

另一方面,如果您需要为每个中间 observable 提供计数,那么您可以将计数移动到 flatMap() 中并打印出计数并在完成时将其重置:

AtomicLong counter = new AtomicLong()
Observable.interval(10, TimeUnit.SECONDS).startWith(0)
        .map { createFinite(ch++) }
        .flatMap { it
                     .doOnNext( counter.incrementAndget()
                     .doOnCompleted( { long ctr = counter.getAndSet(0)
                                        println("I am done. Total event count is $ctr")
                                     } )
        .subscribe { println("Just received [$it] from the infinite stream ") }

这不是很实用,但这种报告往往会破坏正常的流。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2013-06-30
    • 1970-01-01
    • 1970-01-01
    • 2011-09-18
    • 2019-08-27
    • 1970-01-01
    • 2017-12-07
    • 1970-01-01
    相关资源
    最近更新 更多