【问题标题】:RxJava - SwitchMap alike with multiple limited active streamsRxJava - 具有多个有限活动流的 SwitchMap
【发布时间】:2017-06-15 07:56:40
【问题描述】:

我想知道如何以类似于 switchMap 的方式转换可观察对象,但不是限制为单个活动流,而是具有多个(有限)流。

目的是让多个任务同时工作,达到一定的任务计数限制,并允许新任务以先进先出队列策略启动,这意味着任何新任务到达将立即启动,队列中最旧的任务将被取消。

switchMap 将为源的每次发射创建 Observable,并在创建新的 Observable 流后取消之前运行的 Observable 流,我想实现类似但允许某种级别的并发(如 flatMap),这意味着允许创建多个 Observable对于每次发射,并发运行到某个并发限制,当达到并发限制时,最旧的 observable 将被取消,新的 observable 将启动。

其实这也和 flatMap 的 maxConcurrent 类似,只是当达到 maxConcurrent 时,新的 Observables 不会在队列中等待,而是取消旧的 Observables 并立即进入新的 Observables。

【问题讨论】:

  • 我不相信 rxJava(或我熟悉的任何 RX 实现)中存在这样的运算符。如果您有一个好的用例,您可以继续提交功能请求here
  • 嗯,这么想,我会尝试自己想出一些东西,并分享它,看看我是否在正确的方向
  • 这是用于 rxjava 1 还是 2?
  • 它是 rxjava 1 的,已经看过 rxjava-extras Dave,也许我错过了?

标签: concurrency rx-java


【解决方案1】:

虽然没有现成的解决方案,但以下内容应该会有所帮助。

public static void main(String[] args) {

    Observable.create(subscriber -> {
                for (int i = 0; i < 5; i++) {
                    Observable.timer(i, TimeUnit.SECONDS).toBlocking().subscribe();
                    subscriber.onNext(i);
                }
            })
            .switchMap(
                    n -> {
                        System.out.println("Main task emitted event - " + n);
                        return Observable.interval(1, TimeUnit.SECONDS).take((int) n * 3)
                                .doOnUnsubscribe(() -> System.out.println("Unsubscribed for main task event - "+ n));
                    }).subscribe(n2 -> System.out.println("\t" + n2));

    Observable.timer(20, TimeUnit.SECONDS).toBlocking().subscribe();
}

Observable.create 部分创建一个缓慢的生产者,它以发出 0、休眠 1 秒并发出 1、休眠 2 秒并发出 2 等方式发出项目。

switchMap 为每个每秒发出数字的元素创建 Observable 对象。您还可以注意到,每当主 Observable 发出元素时以及取消订阅时,它都会打印一行。

因此,在您的情况下,您可能有兴趣使用doOnUnsubscribe 关闭最旧的任务。希望对您有所帮助。

下面的伪代码可能有助于更好地理解。

getTaskObservable()
        .switchMap(
                task -> {
                    System.out.println("Main task emitted event - " + task);
                    return Observable.create(subscriber -> {
                        initiateTaskAndNotify(task, subscriber);
                    }).doOnUnsubscribe(() -> checkAndKillIfMaxConcurrentTasksReached(task));
                }).subscribe(value -> System.out.println("Done with task and got output" + value));

【讨论】:

  • 谢谢,但我不明白这如何满足我的要求?
  • 我已经用伪代码更新了答案,以帮助更好地理解。请检查并告诉我。
  • 它不能满足我的要求,因为你正在做 switchMap,这将杀死(取消订阅)每个新发射的最后一个订阅者,此外,我需要取消订阅最旧 Observable 的事件是订阅而不是取消订阅.请参阅我在问题中的详细描述。您可以查看 akarnokd 答案代码以了解差异,谢谢。
【解决方案2】:

你可以试试这个转换器:

public static <T, R> Observable.Transformer<T, R> switchFlatMap(
        int n, Func1<T, Observable<R>> mapper) {
    return f -> 
        Observable.defer(() -> {
            final AtomicInteger ingress = new AtomicInteger();
            final Subject<Integer, Integer> cancel = 
                    PublishSubject.<Integer>create().toSerialized();

            return f.flatMap(v -> {
                int id = ingress.getAndIncrement();
                Observable<R> o = mapper.call(v)
                        .takeUntil(cancel.filter(e -> e == id + n));
                cancel.onNext(id);
                return o;
            });
        })
    ;
}

演示:

public static void main(String[] args) {
    PublishSubject<Integer> ps = PublishSubject.create();

    @SuppressWarnings("unchecked")
    PublishSubject<Integer>[] pss = new PublishSubject[3];
    for (int i = 0; i < pss.length; i++) {
        pss[i] = PublishSubject.create();
    }

    AssertableSubscriber<Integer> ts = ps
    .compose(switchFlatMap(2, v -> pss[v]))
    .test();

    ps.onNext(0);
    ps.onNext(1);

    pss[0].onNext(1);
    pss[0].onNext(2);
    pss[0].onNext(3);

    pss[1].onNext(10);
    pss[1].onNext(11);
    pss[1].onNext(12);

    ps.onNext(2);

    pss[0].onNext(4);

    pss[2].onNext(20);
    pss[2].onNext(21);
    pss[2].onNext(22);

    pss[1].onCompleted();
    pss[2].onCompleted();
    ps.onCompleted();

    ts.assertResult(1, 2, 3, 10, 11, 12, 20, 21, 22);
}

【讨论】:

  • 谢谢,这是我的方向,实际上我已经尝试过类似的方法,这种方法的一个问题是 takeUntil 不会取消订阅源 Observable,但只会减少它的排放,所以我们'会有新的 Observable 忽略这些排放。我希望取消订阅源(例如在 switchMap 中),因为我希望释放资源(我的情况是取消网络请求)
  • takeUntil 确实退订了内部资源,您为什么不这么认为?
  • 非常优雅,喜欢!
  • 感谢 akarnokd 的解决方案!!实际上,我已经尝试过使用 takeUntil 的东西,并且因为我期望得到终端事件而感到困惑,因此错误地得出了这个结论。无论如何,现在我可以从包装可观察对象的更复杂的解决方案中退出并返回到 takeUntil 方法,而且您的代码比我提出的要优雅得多,所以谢谢。
猜你喜欢
  • 1970-01-01
  • 2016-07-12
  • 2018-09-04
  • 2015-03-26
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多