【问题标题】:How to properly combine RxJava / Vert.x Data Fetchers / GraphQL Subscriptions如何正确组合 RxJava / Vert.x Data Fetchers / GraphQL Subscriptions
【发布时间】:2019-10-18 14:14:42
【问题描述】:

我有以下测试代码:

  FlowableOnSubscribe<SomeObj> fos;

  private void init() {
    fos = emitter -> {
      try {
        while(true) {
          SomeObj someObj = readFromDataInputStream();
          emitter.onNext(someObj);
          System.out.println("Emitted object");
        }
      }
      catch(Exception e) {
        emitter.onError(e);
      }
    };
  }

  public Single<String> doWork() {
    Flowable<String> myFlow = Flowable.defer(() -> 
      Flowable.create(fos, BackpressureStrategy.BUFFER)
        .cache()
        .subscribeOn(Schedulers.io(), false)
        .doOnSubscribe(x -> doSomethingToTriggerDataInputStream())
        .map(x -> convertMyCustomObjectToString(x))
      );
    );

    // lastOrError/singleOrError end up blocking :(
    return myFlow.firstOrError(); 
  }

  public VertxDataFetcher<CompletionStage<String>> vertxDataFetcherTest() {
    return new VertxDataFetcher<>((env, future) -> {
      try {
        future.complete(doWork().to(SingleInterop.get()));
      }
      catch(Exception e) {
        future.fail();
      }
    });
  }

  public DataFetcher<CompletionStage<String>> dataFetcherTest() {
    return env -> doWork().to(SingleInterop.get());
  }

如果我运行示例代码,它会在第一次成功使用后挂起。换句话说,它会在初始网页加载时运行一次,但如果我在浏览器中进行刷新(Ctrl F5),它会挂起并且不再完成调用。

FWIW,经过一点调试,它看起来像 在 RxCachedThreadScheduler1 上发生 defer/map 调用,在网页刷新后它挂在与 RxCachedThreadScheduler2 的延迟调用中。

如果我切换到为每个呼叫使用单独的输入流,它不会挂断。我遵循的示例在这里(RxJava: Feed one stream (Observable) as the input of another stream...)。

但是,这不适用于我的设计,因为它需要 1 个始终保持打开状态的共享连接。这是因为我希望 GraphQL 订阅能够捕获数据输入流上发出的任何内容。如果我有多个/单独的套接字连接,GraphQL 订阅输入流将错过在非订阅输入流上发出的任何内容。 (除非最好的选择是让所有其他输入流也发送到 GraphQL 订阅流,我不知道该怎么做......)

作为旁注,在这种情况下我使用 VertxDataFetcher 还是 DataFetcher 是否重要?我目前正在使用 DataFetcher 来运行我的示例。如果我需要切换到 VertxDataFetcher,我不确定如何正确转换类型以使该方法正常工作。

【问题讨论】:

    标签: rx-java rx-java2 vert.x


    【解决方案1】:

    发生了很多令人担忧的事情。

    您的发射器明显处于阻塞状态,因此您可以尝试的最重要更改是在创建 Flowable 时调用 .observeOn(RxHelper.blockingScheduler(vertx))。阅读更多here 了解如何使用正确的调度程序。

    不过,这可能无法解决您的问题。您可能正在使用 readFromDataInputStream() 以阻塞的传统 Java IO 方式读取数据,这可能是线程安全的,但几乎可以肯定它不支持并发读取器。只要 readFromDataInputStream() 最终从同一个 IO 读取,您就试图隐式地使用并发读取器。见Java: Concurrent reads on an InputStream

    【讨论】:

    • 您好,感谢您的意见,但我不确定这是我要找的。该示例是针对发射器的一个实例,上面有多个观察者。自从这篇原始帖子以来,我也尝试过 RxHelper,但它没有用(正如你也提到的)
    猜你喜欢
    • 2017-10-24
    • 1970-01-01
    • 2018-09-18
    • 2022-09-29
    • 2020-03-29
    • 1970-01-01
    • 1970-01-01
    • 2022-08-09
    • 1970-01-01
    相关资源
    最近更新 更多