【发布时间】:2017-10-25 03:43:06
【问题描述】:
rxjava2 2.1.5版
试图理解 RxJava2 对 observable 的多个订阅。 有一个简单的文件监视服务来跟踪目录中文件的创建、修改和删除。 我添加了 2 个订阅者,并希望在两个订阅者上都打印事件。 当我将文件复制到监视目录中时,我看到一个订阅者打印出该事件。然后,当我删除文件时,我看到第二个订阅者打印出事件。 我期待两个订阅者都打印事件。我在这里错过了什么?
import static java.nio.file.StandardWatchEventKinds.ENTRY_CREATE;
import static java.nio.file.StandardWatchEventKinds.ENTRY_DELETE;
import static java.nio.file.StandardWatchEventKinds.ENTRY_MODIFY;
import java.io.IOException;
import java.nio.file.FileSystem;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.nio.file.WatchEvent;
import java.nio.file.WatchKey;
import java.nio.file.WatchService;
import java.util.concurrent.TimeUnit;
import io.reactivex.BackpressureStrategy;
import io.reactivex.Flowable;
import io.reactivex.schedulers.Schedulers;
public class MyRxJava2DirWatcher {
public Flowable<WatchEvent<?>> createFlowable(WatchService watcher, Path path) {
return Flowable.create(subscriber -> {
boolean error = false;
WatchKey key;
try {
key = path.register(watcher, ENTRY_CREATE, ENTRY_DELETE, ENTRY_MODIFY);
}
catch (IOException e) {
subscriber.onError(e);
error = true;
}
while (!error) {
key = watcher.take();
for (final WatchEvent<?> event : key.pollEvents()) {
subscriber.onNext(event);
}
key.reset();
}
}, BackpressureStrategy.BUFFER);
}
public static void main(String[] args) throws IOException, InterruptedException {
Path path = Paths.get("c:\\temp\\delete");
final FileSystem fileSystem = path.getFileSystem();
WatchService watcher = fileSystem.newWatchService();
MyRxJava2DirWatcher my = new MyRxJava2DirWatcher();
my.createFlowable(watcher, path).subscribeOn(Schedulers.computation()).subscribe(event -> {
System.out.println("1>>Event kind:" + event.kind() + ". File affected: " + event.context() + ". "
+ Thread.currentThread().getName());
}, onError -> {
System.out.println("1>>" + Thread.currentThread().getName());
onError.printStackTrace();
});
// MyRxJava2DirWatcher my2 = new MyRxJava2DirWatcher();
my.createFlowable(watcher, path).subscribeOn(Schedulers.computation()).subscribe(event -> {
System.out.println("2>>Event kind:" + event.kind() + ". File affected: " + event.context() + ". "
+ Thread.currentThread().getName());
}, onError -> {
System.out.println("2>>" + Thread.currentThread().getName());
onError.printStackTrace();
});
TimeUnit.MINUTES.sleep(1000);
}
}
输出如下所示
2>>Event kind:ENTRY_CREATE. File affected: 1.txt. RxCachedThreadScheduler-2
2>>Event kind:ENTRY_MODIFY. File affected: 1.txt. RxCachedThreadScheduler-2
1>>Event kind:ENTRY_DELETE. File affected: 1.txt. RxCachedThreadScheduler-1
【问题讨论】:
-
您的 WatchService#register 方法是否提供注册多个侦听器的能力,或者每次调用 #register 时都会覆盖该侦听器?如果是这样,那么很明显,第二个订阅会覆盖第一个 #register 侦听器,并且第一个订阅不会再收到通知。只需多播 observable:blog.danlew.net/2016/06/13/multicasting-in-rxjava
-
@hans 当我添加 20 个文件时,我收到 40 个通知,每个订阅者 20 个通知,即每个订阅者 10 个创建 +10 个修改事件。所以它不像第一个订阅者被覆盖。我还尝试创建 flowable 的新实例,但没有任何改变。我将把多播视为我真正想要的,但我还需要了解当前的行为。
标签: rx-java rx-java2 watchservice