0

学习 RxJava 将不胜感激任何建议。

计算准备好后需要接收对象,所以我PublishSubject在models方法中做了:

public PublishSubject<BaseUnit> exec(int inputNumber) {

    if (unitList.size() > 0) {
        for (BaseUnit unit : unitList) {
            unit.setInProgress();
        }
    }

    PublishSubject<BaseUnit> subject = PublishSubject.create();

    list = new ArrayList<>();

    populateList(inputNumber).
            subscribeOn(Schedulers.from(Executors.newFixedThreadPool(ThreadPool.getPoolSize())))
            .subscribe(calculatedList -> {
                list = calculatedList;

                for (List<Integer> elem : list) {
                    for (ListOperationName operationName : ListOperationName.values()) {
                        ListUnit unit = new ListUnit(operationName, elem, 0);
                        calculate(unit);
                        unitList.add(unit);
                        subject.onNext(unit);
                    }
                }

            }, error -> Log.d("ERROR", error.toString()));

    return subject;
}


public Observable<ArrayList<List<Integer>>> populateList(int inputNumber) {
    return Observable.fromCallable(() -> {

        ArrayList<List<Integer>> list = new ArrayList<>();

        Integer[] populatedArray = new Integer[inputNumber];
        Arrays.fill(populatedArray, insertValue);

        list.add(new ArrayList<>(Arrays.asList(populatedArray)));
        list.add(new LinkedList<>(Arrays.asList(populatedArray)));
        list.add(new CopyOnWriteArrayList<>(Arrays.asList(populatedArray)));

        return list;
    });
}

然后尝试订阅演示者:

public void calculate(int inputNumber) {

    fragment.showAllProgressBars();

    repository.exec(inputNumber)
            .observeOn(AndroidSchedulers.mainThread())
            .subscribeOn(Schedulers.from(Executors.newFixedThreadPool(ThreadPool.getPoolSize())))
            .subscribe(unit -> {
                Log.d("PRESENTER RESULT", unit.toString());
                fragment.setCellText(unit.getViewId(), unit.getTimeString());

            }, error -> Log.d("PRESENTER ERROR", error.toString()));
}

这给了我什么。但是如果我使用ReplaySubject- 它会给我所有的结果,但似乎它只使用一个线程。所以我认为我在订阅方面出了点问题,它应该在更早的地方。我需要准确地使用PublishSubject给我的结果,因为它们已经准备好使用多个线程。

如何解决?或者也许还有其他问题?

4

0 回答 0