【问题标题】:RxJava: do authorization parallel with data fetchingRxJava:授权与数据获取并行
【发布时间】:2015-07-29 23:15:52
【问题描述】:

我想结合两个 observables:一个是检查用户是否被授权获取数据,另一个是实际获取数据,我想在获取数据的同时进行授权以减少总延迟。这是一个不平行的例子:

public static void main(String[] args) {
    long startTime = currentTimeMillis();

    Observable<Integer> result = isAuthorized().flatMap(isAuthorized -> {
        if (isAuthorized) return data();
        else throw new RuntimeException("unauthorized");
    });

    List<Integer> data = result.toList().toBlocking().single();
    System.out.println("took: " + (currentTimeMillis() - startTime) + "ms");
    System.out.println(data);
    assert data.size() == 10;
}

private static Observable<Boolean> isAuthorized() {
    return Observable.create( s -> {
        try { sleep(5000); } catch (Exception e) {} // simulate latency
        s.onNext(true);
        s.onCompleted();
    });
}

private static Observable<Integer> data() {
    return Observable.create(s -> {
        for (int i = 0; i < 10; i++) {
            try { sleep(1000); } catch (Exception e) {} // simulate long running job
            s.onNext(i);
        }
        s.onCompleted();
    });
}

执行此操作的总时间为 15 秒,如果对授权和数据获取的调用是并行的,则应为 10 秒。怎么做?理想情况下,我还想在等待授权完成时最多缓存多少数据项在内存中。

顺便说一句,我已经阅读了excellent answer about paralleling observables,但现在仍然不知道如何解决我的问题。

【问题讨论】:

    标签: java reactive-programming rx-java


    【解决方案1】:

    要在类型安全的情况下做到这一点,我建议为isAuthorized() 排放和data() 排放使用包装不可变类并合并流,然后减少和过滤以不发射任何内容(未经授权)或数据(授权)。

    static class AuthorizedData {
    
        final Boolean isAuthorized; //null equals unknown
        final Data data; //null equals unknown
    
        AuthOrData(Boolean isAuthorized, Data data) {
            this.isAuthorized = isAuthorized;
            this.data = data;
        }
    
    }
    
    Observable<Data> authorizedData =
      isAuthorized()
        .map(x -> new AuthorizedData(x, null))
        .subscribeOn(Schedulers.io())
        .mergeWith(
            data().map(x -> new AuthorizedData(null, x))
                  .subscribeOn(Schedulers.io()))
        .takeUntil(a -> a.isAuthorized!=null && !a.isAuthorized)
        .reduce(new AuthorizedData(null, null), (a, b) -> {
            if (a.isAuthorized!=null && a.data != null)
               return a;
            else if (b.isAuthorized!=null)
               return new AuthorizedData(b.isAuthorized, a.data);
            else if (b.data!=null)
               return new AuthorizedData(a.isAuthorized, b.data);
            else 
               return a;
        })
        .filter(a -> a.isAuthorized!=null 
                     && a.isAuthorized && a.data!=null)
        .map(a -> a.data);
    

    上面的authorizedData如果未授权则为空,否则为单个数据项的流。

    上面takeUntil的意思是当发现用户没有被授权时,立即退订data()。如果data() 是可中断的(可以关闭套接字或其他),这将非常有用。

    【讨论】:

    • 更新了 takeUntil 添加取消订阅 data() 如果未经授权
    • 此解决方案有效,但前提是 data() 返回单个项目。我的情况是它返回了许多项目,我必须使用 isAuthorized() 调用来保护所有项目,该调用返回单个布尔值,表示用户是否可以看到所有数据项。
    【解决方案2】:

    我已经找到了一种方法:

    1. 使用 Schedulers.io() 订阅两个 observable 以并行运行它们。
    2. 使用repeat() 无限次重复isAuthorized() 发射,因为它通常只发射单个布尔值。
    3. 为避免为每个项目调用授权服务,请使用cache()
    4. 重复布尔值与项目流压缩流,并使用布尔值决定是返回项目还是抛出异常。

    解决办法如下:

    public static void main(String[] args) {
      long startTime = currentTimeMillis();
    
      Observable<Integer> result = Observable.defer(() -> {
        Observable<Boolean> p1 = isAuthorized().cache().repeat().subscribeOn(Schedulers.io());
        Observable<Integer> p2 = data().subscribeOn(Schedulers.io());
        return Observable.zip(p1, p2, (isAuthorized, item) -> {
          if (isAuthorized)
            return item;
          else
            throw new RuntimeException("unauthorized");
        });
      });
    
      List<Integer> data = result.toList().toBlocking().single();
      System.out.println("took: " + (currentTimeMillis() - startTime) + "ms");
      System.out.println(data);
      assert data.size() == 10;
    }
    
    private static Observable<Boolean> isAuthorized() {
      return Observable.create( s -> {
        try { sleep(5000); } catch (Exception e) {} // simulate latency
        s.onNext(true);
        s.onCompleted();
      });
    }
    
    private static Observable<Integer> data() {
      return Observable.range(1, 10)
          .doOnNext(i -> { try { sleep(1000); } catch (Exception e) {} });
    }
    

    据我观察,此解决方案还避免了 OutOfMemory 错误。两个可观察对象同时开始发出,但如果授权服务较慢,数据项将被收集,直到内部缓冲区被填满。然后 RxJava 将停止请求数据项,直到授权最终发出布尔值。

    当授权返回否定结果时,它也会正确地从数据流中取消订阅。

    【讨论】:

    • 好主意。但是如何处理未经授权的异常。你很高兴收到 onError 吗?如果使用自定义异常类型和。 .doOnErrorResumeNext 你也可以抑制它。
    • 我们实际上使用异常来控制流程:在 servlet 级别,我们订阅安全的 observable,如果自定义授权异常出错,我们返回 HTTP 403 Forbidden,否则我们将每个数据项序列化到 servlet 输出流。
    猜你喜欢
    • 2014-12-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-11-21
    • 2020-12-28
    • 1970-01-01
    相关资源
    最近更新 更多