【问题标题】:RxJava: Split Rx Flowable into multiple streamsRxJava:将 Rx Flowable 拆分为多个流
【发布时间】:2017-05-05 03:59:33
【问题描述】:

我想对stream进行一些操作,然后将stream分成两个stream,然后分别处理。

显示问题的示例:

Flowable<SuccessfulObject> stream = Flowable.fromArray(
        new SuccessfulObject(true, 0),
        new SuccessfulObject(false, 1),
        new SuccessfulObject(true, 2));

stream = stream.doOnEach(System.out::println);

Flowable<SuccessfulObject> successful = stream.filter(SuccessfulObject::isSuccess);
Flowable<SuccessfulObject> failed = stream.filter(SuccessfulObject::isFail);

successful.doOnEach(successfulObject -> {/*handle success*/}).subscribe();
failed.doOnEach(successfulObject -> {/*handle fail*/}).subscribe();

类:

class SuccessfulObject {
    private boolean success;
    private int id;

    public SuccessfulObject(boolean success, int id) {
        this.success = success;
        this.id = id;
    }

    public boolean isSuccess() {
        return success;
    }
    public boolean isFail() {
        return !success;
    }

    public void setSuccess(boolean success) {
        this.success = success;
    }

    @Override
    public String toString() {
        return "SuccessfulObject{" +
                "id=" + id +
                '}';
    }
}

但是这段代码将所有元素打印两次,而我想在只拆分一次之前执行所有操作。

输出:

OnNextNotification[SuccessfulObject{id=0}]
OnNextNotification[SuccessfulObject{id=1}]
OnNextNotification[SuccessfulObject{id=2}]
OnCompleteNotification
OnNextNotification[SuccessfulObject{id=0}]
OnNextNotification[SuccessfulObject{id=1}]
OnNextNotification[SuccessfulObject{id=2}]
OnCompleteNotification

如何处理流以接收此行为?

【问题讨论】:

  • 是否要将处理结果合并回一个流(fork-join-behaviour?)
  • 不,只是拆分流并分别执行所有操作。
  • 好吧,然后使用@akarnokd 的解决方案。作为侧节点:不要在 rx-pipeline 中使用可变对象。 isFail 也不是必需的,因为 isSuccess 在 fals 上暗示它失败了。

标签: java stream rx-java


【解决方案1】:

使用publish 分享对源的订阅:

Flowable<Integer> source = Flowable.range(1, 5);

ConnectableFlowable<Integer> cf = source.publish();

cf.filter(v -> v % 2 == 0).subscribe(v -> System.out.println("Even: " + v));

cf.filter(v -> v % 2 != 0).subscribe(v -> System.out.println("Odd: " + v));

cf.connect();

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-12-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多