代码之家  ›  专栏  ›  技术社区  ›  ESala

多线程RxJava的反应式拉动

  •  3
  • ESala  · 技术社区  · 11 年前

    我正在尝试建立一个 反作用力拉力 观察员 RxJava语言 .

    我的观察者是这样的:

    Observable<Command> myObs = Observable.create(s -> {
       Command command;
       int i = 0;
       do {
          command = NetworkOperation1.call(i);
          logger.info("Init command " + i);
          s.onNext(command);
          i++;
       } while (!command.isLast() && i < MAX);
       s.onCompleted();
    });
    

    我想在4个并发批处理(缓冲区)中处理它,如下所示:

    myObs
        .buffer(10)
        .flatMap(batch -> {
              return Observable
                       .from(batch)
                       .subscribeOn(Schedulers.io())
                       .map(c -> {
                           Intermediate m = NetworkOperation2.call(c));
                           logger.info("Done intermediate " + m.id);
                           return m;
                       }
              }, 4);
    

    然后,我需要以不同的大小批量处理结果,如下所示:

        .buffer(25)
        .subscribeOn(Schedulers.newThread())
        .subscribe(list ->
             logger.info("Finished batch with " + list.size());
    

    问题是Observable中的Commands 一次处理所有 ,而我希望处理它们 根据需要 .

    下面是发生的情况的日志:(请注意,所有1000个命令都会同时运行,而不是根据需要调用)

    Init command 0
    Init command 1
    Init command 2
    ...
    Init command 999
    Done intermediate 0
    Done intermediate 1
    ...
    Done intermediate 24
    Finished batch with 25
    Done intermediate 25
    Done intermediate 26
    ...
    Done intermediate 49
    Finished batch with 25
    ...
    

    问题 :有没有一种方法可以暂停Observer的线程,这样它就不会同时执行所有命令或类似的操作?我尝试了request()运算符,但无法使其工作。

    非常感谢。

    1 回复  |  直到 11 年前
        1
  •  4
  •   Dave Moten    11 年前

    您需要背压感知源和操作员。您使用的运算符支持背压,但您的源不支持。

    请改为执行以下操作:

    myObs = Observable.range(1,1000)
        .map(i -> NetworkOperation1.call(i));
    

    Observable.range 支持背压,因此只有在请求时才会发出。