【发布时间】: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<List<String>> collector?
2.为了按顺序应用转换,我使用了一个 for 循环;
我尝试了多种变体(即:flatMap、zipWith),但是,最终发生的是转换没有按顺序应用;如何在没有 for 循环的情况下对此进行建模?
for (Transform t : transforms) {
stage = stage.flatMap(t::applyAsync);
}
基本上,我需要一种在输入 List<List<String>> 上应用 Observable<List<List<String>>> applyAsync(List<List<String>> input) 并在每次转换 (Observable<List<List<String>>>) 上递归地继续这样做的方法。
类似于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