【问题标题】:RxJava2, 2 subscribers for an Observable/Flowable but onNext getting called on any oneRxJava2,一个 Observable/Flowable 的 2 个订阅者,但 onNext 被任何一个订阅者调用
【发布时间】: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


【解决方案1】:

发生的情况是,您在两个 Flowables 之间共享相同的 WatchService,并且他们在其中竞争事件。如果您改为传入FileSystem 并在Flowable.create 中调用newWatchService(),您应该会收到与Subscribers 一样多的所有事件:

public Flowable<WatchEvent<?>> createFlowable(FileSystem fs, Path path) {

    return Flowable.create(subscriber -> {

        WatchService watcher = fs.newWatchService();

        subscriber.setCancellable(() -> watcher.close());

        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);

}

还请注意,您应该使用subscribeOn(Schedulers.computation(), false) 以避免poll 与您的Subscriber 死锁。

【讨论】:

  • 现在按预期工作。谢谢。
  • 我不理解推荐的 javadoc 或错误标志。如果在 * 链中有一个 {@link #create(FlowableOnSubscribe, BackpressureStrategy)} 类型的源,建议将 {@code requestOn} 设置为 false 以避免相同池死锁 * 因为请求可能会堆积在一个 Eager/阻塞发射器
  • Here 是关于它的讨论。
【解决方案2】:

您正在为两个不同的订阅者创建两个不同的 Flowable。是否有一个 Flowable 被订阅两次,如下所示。

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();
        Flowable myFlowable = my.createFlowable(watcher, path);

        myFlowable.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();
        });

        myFlowable.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);

    }
}

【讨论】:

  • 这对我没有任何改变。结果与我的代码相同。只是输出的一个 sn-p.... 我没有得到线程池 1 和线程池 2 1>>事件类型:ENTRY_CREATE 中每个文件的创建和修改事件。受影响的文件:1.txt。 RxComputationThreadPool-1 1>>事件类型:ENTRY_MODIFY。受影响的文件:1.txt。 RxComputationThreadPool-1 2>>事件类型:ENTRY_CREATE。受影响的文件:2 - 复制 (10).txt。 RxComputationThreadPool-2 1>>事件类型:ENTRY_MODIFY。受影响的文件:2 - 复制 (10).txt。 RxComputationThreadPool-1
  • 我稍后会删除我的答案,因为@akarnokd 在 Rx 方面是王者
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-12-03
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-09-15
  • 2017-01-30
相关资源
最近更新 更多