【问题标题】:Looking for better ways to compose RxJava Observables寻找更好的方法来编写 RxJava Observables
【发布时间】:2017-06-20 12:31:11
【问题描述】:

我正在开发一种从多个来源收集数据并按顺序应用一些转换的工具。我目前正在将此功能从 Java 8 流转换为使用 ReactiveX/RxJava。

您可以在下面看到一个演示当前 RxJava 实现的单元测试。

虽然它有效,但我对结果不够满意,正在寻找有关如何改进它的指导!


我的两个问题是:

1.每个源返回一个结果列表 (List>)。因为需要对整个数据集进行转换,所以我需要将多个列表合并为一个。

现在代码如下所示:

Observable<List<List<String>>> stage = Observable.merge(src1, src2, src3, src4);

final List<List<String>> collector = new ArrayList<>();
Single<List<List<String>>> combinedData = stage.reduce(collector, (list, items) -> {
    list.addAll(items);
    return list;
});

有没有办法摆脱可观察流之外的List&lt;List&lt;String&gt;&gt; collector?


2.为了按顺序应用转换,我使用了一个 for 循环; 我尝试了多种变体(即:flatMap、zipWith),但是,最终发生的是转换没有按顺序应用;如何在没有 for 循环的情况下对此进行建模?

for (Transform t : transforms) {
    stage = stage.flatMap(t::applyAsync);
}

基本上,我需要一种在输入 List&lt;List&lt;String&gt;&gt; 上应用 Observable&lt;List&lt;List&lt;String&gt;&gt;&gt; applyAsync(List&lt;List&lt;String&gt;&gt; input) 并在每次转换 (Observable&lt;List&lt;List&lt;String&gt;&gt;&gt;) 上递归地继续这样做的方法。

类似于Observable.reduce,但累加器函数需要在每次迭代时改变。


这是我写的完整的单元测试代码:

import io.reactivex.Observable;
import io.reactivex.schedulers.Schedulers;
import org.mockito.ArgumentCaptor;
import org.testng.annotations.Test;

import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;

import static org.hamcrest.MatcherAssert.assertThat;
import static org.hamcrest.Matchers.contains;
import static org.hamcrest.Matchers.*;
import static org.mockito.Mockito.*;

public class ObservableTest {
    @Test
    public void testObservable() throws Exception {
        // ARRANGE
        final CountDownLatch DONE = new CountDownLatch(1);

        // init source objects
        Observable<List<List<String>>> src1 = Observable.just(makeMatrix(Arrays.asList("src11")));
        Observable<List<List<String>>> src2 = Observable.just(makeMatrix(Arrays.asList("src21", "src22")));
        Observable<List<List<String>>> src3 = Observable.just(makeMatrix(Arrays.asList("src31", "src32", "src33")));
        Observable<List<List<String>>> src4 = Observable.just(makeMatrix(Arrays.asList("src41"), Arrays.asList("src51")));

        // prepare transformations and processor
        List<Transform> transforms = Arrays.asList(new Transform(1, 100), new Transform(2, 0));
        Processor processor = spy(new Processor());


        // ACT

        // Concat sources
        Observable<List<List<String>>> stage = Observable.merge(src1, src2, src3, src4);

        // Merge individual into matrix

        // (#1) Can the reduce operation be written without the accumulator?
        final List<List<String>> collector = new ArrayList<>();
        Single<List<List<String>>> combinedData = stage.reduce(collector, (list, items) -> {
            list.addAll(items);
            return list;
        });


        // Transform
        stage = combinedData.toObservable();
        for (Transform t : transforms) {
            // (#2) Can a series of transforms be applied sequentially to a Single (List<List<String>>), without the use of a for-loop?
            stage = stage.flatMap(t::applyAsync);
        }

        // Process
        stage.doOnComplete(DONE::countDown)
                .subscribeOn(Schedulers.computation())
                .subscribe(o -> System.out.println(processor.printList(o)));

        // wait for processing to complete
        DONE.await();


        // ASSERT

        // The sources should be combined in a single matrix
        @SuppressWarnings("unchecked")
        ArgumentCaptor<List<List<String>>> resultCaptor = ArgumentCaptor.forClass(List.class);

        verify(processor, times(1)).printList(resultCaptor.capture());
        List<List<String>> resultMatrix = resultCaptor.getValue();

        // result matrix should not be null and all transformations should be applied in order (T1, T2, etc.)
        assertThat(resultMatrix, notNullValue());
        assertThat(resultMatrix.stream().flatMap(Collection::stream).collect(Collectors.toList()), everyItem(containsString("T1-T2")));
        assertThat(resultMatrix, not(hasItem(hasItem(containsString("T2-T1")))));
   }


    private List<List<String>> makeMatrix(List<String> items) {
        return Collections.singletonList(items);
    }

    private List<List<String>> makeMatrix(List<String> items, List<String> moreItems) {
        return Arrays.asList(items, moreItems);
    }

    static class Processor {
        String printList(List<List<String>> input) {
            return input.stream().map(rows -> rows.stream().collect(Collectors.joining(" | ")))
                    .collect(Collectors.joining(System.lineSeparator()));
        }
    }

    static class Transform {
        final int n;
        private final int delay;

        Transform(int n, int delay) {
            this.n = n;
            this.delay = delay;
        }

        private Observable<List<List<String>>> applyAsync(List<List<String>> input) {
            return Observable.just(input).map(this::apply).delay(delay, TimeUnit.MILLISECONDS);
        }

        private List<List<String>> apply(List<List<String>> input) {
            return input.stream()
                    .map(row -> row.stream()
                            .map(this::transform)
                            .collect(Collectors.toList())
                    )
                    .collect(Collectors.toList());
        }

        private String transform(String input) {
            return input + "-T" + n;
        }
    }
}

如果要运行,请导入以下 Maven 依赖项:

<dependency>
    <groupId>io.reactivex.rxjava2</groupId>
    <artifactId>rxjava</artifactId>
    <version>2.1.0</version>
</dependency>
<dependency>
    <groupId>org.hamcrest</groupId>
    <artifactId>hamcrest-all</artifactId>
    <version>1.3</version>
    <scope>test</scope>
</dependency>
<dependency>
    <groupId>org.mockito</groupId>
    <artifactId>mockito-core</artifactId>
    <version>2.4.3</version>
    <scope>test</scope>
</dependency>
<dependency>
    <groupId>org.testng</groupId>
    <artifactId>testng</artifactId>
    <version>6.10</version>
    <scope>test</scope>
</dependency>

【问题讨论】:

  • flatMap 如果与并发一起使用,则不会(必然)保留顺序。
  • 如果您需要更改函数,为什么不让您的累加器函数调用您返回的某个函数。如果您需要返回一个函数并希望发出其他内容,您可以发出一个Pair,然后下一个运算符可以是一个映射,您可以将它映射到您想要的实际发射(不是Function )

标签: java observable reactive-programming rx-java2 reactivex


【解决方案1】:

免责声明:我没有任何 RxJava 特定经验,只有 RxJS、Rx.NET 和 RxSwift

1

您应该可以直接传递ArrayList 的新实例:

Single<List<List<String>>> combinedData = stage.reduce(new ArrayList<>(), (list, items) -> {
    list.addAll(items);
    return list;
});

它只会用作累加器的初始种子;在第一次调用 lambda 之后,它就不再相关了。

2

我认为您希望对这个问题进行一些递归。以下解决方案可能适合您:

// Put this somewhere
public IObservable<List<List<String>> handleTransforms(
    Observable<List<List<String>>> currentStage
    List<Transform> ts){

    return currentStage.flatMap(t[0]::applyAsync)
        .flatMap(newStage -> handleTransforms(newStage, ts.stream().skip(1).toList()))
}

// And then use it like this
stage = handleTransforms(stage, transforms);

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2022-01-13
    • 2020-09-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-04-04
    相关资源
    最近更新 更多