【问题标题】:RxJava: How to .zip two Observable, then .merge them and eventually .reduce to aggregate all resultsRxJava:如何 .zip 两个 Observable,然后 .merge 它们并最终 .reduce 聚合所有结果
【发布时间】:2015-05-04 12:28:27
【问题描述】:

我有以下代码:

public void foo() {
    Long[] gData = new Long[] { 1L, 2L };

    rx.Observable.from(gData)
    .concatMap(data -> {

        rx.Observable<GmObject> depositObs1 = depositToUserBalance(data, 1);
        rx.Observable<GmObject> depositObs2 = depositToUserBalance(data, 2);


        return rx.Observable.zip(depositObs1, depositObs2, (depositObj1, depositObj2) -> {

            depositObj1.putNumber("seat_index", data);
            depositObj2.putNumber("seat_index", data);

            return rx.Observable.merge(
                    rx.Observable.just(depositObj1),
                    rx.Observable.just(depositObj2));
        })
    })
    .reduce(new ArrayList<Long>(), (payoutArr, payoutObj) -> {

        int seatIndex = ((GmObject) payoutObj).getNumber("seat_index").intValue();
        long payout = ((GmObject) payoutObj).getNumber("payout").longValue();
        payoutArr.add(seatIndex, payout);
        return payoutArr;
    })
    .subscribe(results -> {
        System.out.println(results);
    });
}

此代码使用 .zip 向 observables 发出数据,然后添加一个“seat_index”属性并调用 .merge 以使用 .reduce,因此最终所有结果都将聚合到一个 ArrayList 中。

此代码存在问题:当 .reduce 处理其输入时,它会将其作为 Observable 而不是 GmObject ...什么函数可以从其 Observable 包装中“提取” GmObject?

以这种方式使用 rxJava 有意义吗?还是有更好的技术?

谢谢!

【问题讨论】:

    标签: java rx-java


    【解决方案1】:

    zip 运算符将 lambda 作为第三个参数。这个 lambda 是一个 2 args 函数,它返回一个对象,该对象是 args 组合的结果。而不是组合结果的Observable(当然,对象可以是Observable,但这不是您想要的)。

    因此,在您拨打zip 之后,您将获得一个Observable&lt;Observable&lt;GmObject&gt;&gt;,但您希望获得一个Observable&lt;GmObject&gt;

    我认为zip 运算符不是您要查找的运算符。

    public void foo() {
        Long[] gData = new Long[] { 1L, 2L };
    
        rx.Observable.from(gData)
        .concatMap(data -> {
    
            rx.Observable<GmObject> depositObs1 = depositToUserBalance(data, 1).doOnNext(obj -> obj.putNumber("seat_index", data));
            rx.Observable<GmObject> depositObs2 = depositToUserBalance(data, 2).doOnNext(obj -> obj.putNumber("seat_index", data));
    
    
            return rx.Observable.merge(depositObs1, depositObs2);
        })
       .reduce(new ArrayList<Long>(), (payoutArr, payoutObj) -> {
    
            int seatIndex = ((GmObject) payoutObj).getNumber("seat_index").intValue();
            long payout = ((GmObject) payoutObj).getNumber("payout").longValue();
            payoutArr.add(seatIndex, payout);
            return payoutArr;
       })
       .subscribe(results -> System.out.println(results));
    }
    

    【讨论】:

    • 这么简单,但我自己也弄不明白……我从来没有用过doOnNext()……
    • @dwursteisen 有什么方法可以在 Android 上做到这一点(没有 lambda)?
    • 用匿名类替换 lambda。 (尝试 Java 8 项目并要求您的 IDE 将其转换为 Anonymous 类以查看我的意思)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-08-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-01-25
    • 2016-07-31
    相关资源
    最近更新 更多