【问题标题】:Java - processing documents in parallelJava - 并行处理文档
【发布时间】:2016-03-02 09:27:52
【问题描述】:
我有 5 个文档(比如说),我对每个文档都有一些处理。此处的处理包括打开文档/文件、读取数据、进行一些文档操作(编辑文本等)。对于文档操作,我可能会使用 docx4j 或 apache-poi。但我的用例是这样的——我想以某种方式并行处理这 4-5 个文档,利用我 CPU 上可用的多个内核。对每个文档的处理是相互独立的。
在 Java 中实现这种并行处理的最佳方式是什么。我之前在 java 中使用过 ExecutorService 和 Thread 类。但我对Streams 或RxJava 等较新的概念不太了解。可以通过使用 Java 8 中介绍的 Java 中的 Parallel Stream 来完成此任务吗?使用 Executors/Streams/Thread 类等会更好。如果可以使用 Streams,请提供一个链接,我可以在其中找到一些关于如何做到这一点的教程。感谢您的帮助!
【问题讨论】:
标签:
java
concurrency
parallel-processing
java-stream
【解决方案1】:
您可以使用 Java Streams 使用以下模式进行并行处理。
List<File> files = ...
files.parallelStream().forEach(f -> process(f));
或
File[] files = dir.listFiles();
Stream.of(files).parallel().forEach(f -> process(f));
注意:process 在此示例中不能抛出 CheckedException。我建议你要么记录它,要么返回一个结果对象。
【解决方案2】:
如果你想了解 ReactiveX,我推荐使用 rxJava Observable.zip http://reactivex.io/documentation/operators/zip.html
您可以在这里并行运行多个进程的示例:
public class ObservableZip {
private Scheduler scheduler;
private Scheduler scheduler1;
private Scheduler scheduler2;
@Test
public void testAsyncZip() {
scheduler = Schedulers.newThread();//Thread to open and read 1 file
scheduler1 = Schedulers.newThread();//Thread to open and read 1 file
scheduler2 = Schedulers.newThread();//Thread to open and read 1 file
Observable.zip(obAsyncString(file1), obAsyncString1(file2), obAsyncString2(file3), (s, s2, s3) -> s.concat(s2)
.concat(s3))
.subscribe(result -> showResult("All files in one:", result));
}
public void showResult(String transactionType, String result) {
System.out.println(result + " " +
transactionType);
}
public Observable<String> obAsyncString(File file) {
return Observable.just(file)
.observeOn(scheduler)
.doOnNext(val -> {
//Here you read your file
});
}
public Observable<String> obAsyncString1(File file) {
return Observable.just(file)
.observeOn(scheduler1)
.doOnNext(val -> {
//Here you read your file 2
});
}
public Observable<String> obAsyncString2(File file) {
return Observable.just(file)
.observeOn(scheduler2)
.doOnNext(val -> {
//Here you read your file 3
});
}
}
就像我说的,以防万一你想了解 ReactiveX,因为如果不是,在你的堆栈中添加这个框架来解决这个问题会有点矫枉过正,我更喜欢以前的流并行解决方案