【问题标题】:Generate infinite sequence of Natural numbers using RxJava使用 RxJava 生成无限的自然数序列
【发布时间】:2015-07-14 18:47:20
【问题描述】:

我正在尝试使用 RxJava 编写一个简单的程序来生成无限的自然数序列。所以,到目前为止,我已经找到了两种使用 Observable.timer()Observable.interval() 生成数字序列的方法。我不确定这些功能是否是解决此问题的正确方法。我期待一个像我们在 Java 8 中那样的简单函数来生成无限自然数。

IntStream.iterate(1, value -> value +1).forEach(System.out::println);

我尝试将 IntStream 与 Observable 一起使用,但无法正常工作。它仅向第一个订阅者发送无限的数字流。如何正确生成无限自然数列?

import rx.Observable;
import rx.functions.Action1;

import java.util.stream.IntStream;

public class NaturalNumbers {

    public static void main(String[] args) {
        Observable<Integer> naturalNumbers = Observable.<Integer>create(subscriber -> {
            IntStream stream = IntStream.iterate(1, val -> val + 1);
            stream.forEach(naturalNumber -> subscriber.onNext(naturalNumber));
        });

        Action1<Integer> first = naturalNumber -> System.out.println("First got " + naturalNumber);
        Action1<Integer> second = naturalNumber -> System.out.println("Second got " + naturalNumber);
        Action1<Integer> third = naturalNumber -> System.out.println("Third got " + naturalNumber);
        naturalNumbers.subscribe(first);
        naturalNumbers.subscribe(second);
        naturalNumbers.subscribe(third);

    }
}

【问题讨论】:

    标签: java system.reactive reactive-programming rx-java


    【解决方案1】:

    问题是naturalNumbers.subscribe(first);,您实现的OnSubscribe 正在被调用,并且您正在无限流上执行forEach,因此您的程序永远不会终止。

    您可以处理它的一种方法是在不同的线程上异步订阅它们。为了轻松查看结果,我不得不在 Stream 处理中引入睡眠:

    Observable<Integer> naturalNumbers = Observable.<Integer>create(subscriber -> {
        IntStream stream = IntStream.iterate(1, i -> i + 1);
        stream.peek(i -> {
            try {
                // Added to visibly see printing
                Thread.sleep(50);
            } catch (InterruptedException e) {
            }
        }).forEach(subscriber::onNext);
    });
    
    final Subscription subscribe1 = naturalNumbers
        .subscribeOn(Schedulers.newThread())
        .subscribe(first);
    final Subscription subscribe2 = naturalNumbers
        .subscribeOn(Schedulers.newThread())
        .subscribe(second);
    final Subscription subscribe3 = naturalNumbers
        .subscribeOn(Schedulers.newThread())
        .subscribe(third);
    
    Thread.sleep(1000);
    
    System.out.println("Unsubscribing");
    subscribe1.unsubscribe();
    subscribe2.unsubscribe();
    subscribe3.unsubscribe();
    Thread.sleep(1000);
    System.out.println("Stopping");
    

    【讨论】:

    • 感谢迈克的回答。如果我在创建 Observable 时调用 subscribeOn 方法而不是如上面的代码 sn-p 所示调用它三次,会不会有什么不同。我对其进行了测试,行为相同,但仍想确认。
    • 这个问题已被正确识别,但这是个糟糕的建议 - 您永远不应该使用 subscribeOn 来解决此问题 - 请参阅我的答案了解原因。
    • 以这种方式调用unsubscribe 会断开订阅者的连接,因此它会停止接收消息,但它不会停止生成器的循环,它会继续运行无限地消耗您的 CPU 功率。请参阅我的回答,了解如何解决故事的两面。
    【解决方案2】:

    Observable.Generate 正是响应式解决这类问题的算子。我还假设这是一个教学示例,因为为此使用可迭代对象可能会更好。

    您的代码在订阅者的线程上生成整个流。由于它是一个无限流,subscribe 调用永远不会完成。除了这个明显的问题之外,取消订阅也会有问题,因为您没有在循环中检查它。

    您想使用调度程序来解决这个问题——当然不要使用subscribeOn,因为这会给所有观察者带来负担。安排每个号码的发送到onNext - 作为每个预定操作的最后一步,安排下一个。

    基本上这就是Observable.generate 给你的——每次迭代都在提供的调度器上调度(如果你不指定它,默认为引入并发的一个)。可以取消调度程序操作并避免线程饥饿。

    Rx.NET 是这样解决的(实际上有一个更好的 async/await 模型,但在 Java afaik 中不可用):

    static IObservable<int> Range(int start, int count, IScheduler scheduler)
    {
        return Observable.Create<int>(observer =>
        {
            return scheduler.Schedule(0, (i, self) =>
            {
                if (i < count)
                {
                    Console.WriteLine("Iteration {0}", i);
                    observer.OnNext(start + i);
                    self(i + 1);
                }
                else
                {
                    observer.OnCompleted();
                }
            });
       });
    }
    

    这里需要注意两点:

    • 对 Schedule 的调用会返回一个订阅句柄,该句柄会传回观察者
    • 调度是递归的 - self 参数是对用于调用下一次迭代的调度程序的引用。这允许取消订阅以取消操作。

    不确定这在 RxJava 中是怎样的,但想法应该是一样的。同样,Observable.generate 对您来说可能会更简单,因为它旨在处理这种情况。

    【讨论】:

      【解决方案3】:

      创建无限序列时应注意:

      1. 订阅和观察不同的线程;否则您将只为单个订阅者提供服务
      2. 订阅终止后立即停止生成值;否则失控的循环会吃掉你的 CPU

      第一个问题通过使用subscribeOn()observeOn()和各种调度器解决。

      第二个问题最好使用库提供的方法Observable.generate()Observable.fromIterable() 来解决。他们会进行适当的检查。

      检查一下:

      Observable<Integer> naturalNumbers =
              Observable.<Integer, Integer>generate(() -> 1, (s, g) -> {
                  logger.info("generating {}", s);
                  g.onNext(s);
                  return s + 1;
              }).subscribeOn(Schedulers.newThread());
      Disposable sub1 = naturalNumbers
              .subscribe(v -> logger.info("1 got {}", v));
      Disposable sub2 = naturalNumbers
              .subscribe(v -> logger.info("2 got {}", v));
      Disposable sub3 = naturalNumbers
              .subscribe(v -> logger.info("3 got {}", v));
      
      Thread.sleep(100);
      
      logger.info("unsubscribing...");
      sub1.dispose();
      sub2.dispose();
      sub3.dispose();
      Thread.sleep(1000);
      
      logger.info("done");
      

      【讨论】:

        猜你喜欢
        • 2014-12-04
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2016-10-12
        • 1970-01-01
        • 1970-01-01
        • 2018-12-01
        相关资源
        最近更新 更多