【问题标题】:Custom Observable stream of REST results from a list of links来自链接列表的自定义 Observable REST 结果流
【发布时间】:2015-06-10 06:57:17
【问题描述】:

我正在使用一个实现自己的 REST 请求的库。它需要一个“永久链接”并发出一个 REST 请求来获取内容(我在源代码中提供了模拟以进行测试)。我有一个永久链接列表,并且想为每个链接安排一个请求并生成内容结果流。所有这些都应该是异步的,一旦所有结果都完成了,我想发出一个它们的列表。

我正在尝试使用 RxJava 来实现这一点。这是我现在所拥有的(抱歉没有使用 lambda,我只是想习惯 RxJava 类名):

public class Main {

    public static void main(String[] args) {
        int count = 10;
        List<String> permalinks = new ArrayList<>(count);
        for (int i = 1; i <= count; ++i) {
            permalinks.add("permalink_" + i);
        }

        ContentManager cm = new ContentManager();

        Observable
                .create(new Observable.OnSubscribe<Content>() {

                    int gotCount = 0;

                    @Override
                    public void call(Subscriber<? super Content> subscriber) {
                        for (String permalink : permalinks) { // 1. is iterating here the correct way?
                            if (!subscriber.isUnsubscribed()) { // 2. how often and where should I check isUnsubscribed?
                                cm.getBasicContentByPermalink(permalink, new RestCallback() {

                                    @Override
                                    public void onSuccess(Content content) {
                                        if (!subscriber.isUnsubscribed()) {
                                            subscriber.onNext(content);
                                            completeIfFinished(); // 3. if guarded by isUnsubscribed, onComplete might never be called
                                        }
                                    }

                                    @Override
                                    public void onFailure(int code, String message) {
                                        if (!subscriber.isUnsubscribed()) {
                                            subscriber.onNext(null); // 4. is this OK or is there some other way to mark a failure?
                                            completeIfFinished();
                                        }
                                    }

                                    private void completeIfFinished() {
                                        ++gotCount;
                                        if (gotCount == permalinks.size()) { // 5. how to know that the last request is done? am I supposed to implement such custom logic?
                                            subscriber.onCompleted();
                                        }
                                    }
                                });
                            }
                        }
                    }
                })
                .toList()
                .subscribe(new Action1<List<Content>>() {

                    @Override
                    public void call(List<Content> contents) {
                        System.out.println("list count: " + contents.size());
                        System.out.println("results: ");
                        contents.stream().map(content -> content != null ? content.basic : null).forEach(System.out::println);
                    }
                });

        System.out.println("finishing main");
    }
}

// library mocks

class Content {

    String basic;

    String extended;

    public Content(String basic, String extended) {
        this.basic = basic;
        this.extended = extended;
    }
}

interface RestCallback {

    void onSuccess(Content content);

    void onFailure(int code, String message);
}

class ContentManager {

    private final Random random = new Random();

    public void getBasicContentByPermalink(String permalink, RestCallback callback) {
        // just to simulate network latency and unordered results
        new Thread() {

            @Override
            public void run() {
                try {
                    Thread.sleep(random.nextInt(1000) + 200);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                if (random.nextInt(100) < 95) {
                    // 95% of the time we succeed
                    callback.onSuccess(new Content(permalink + "_basic", null));
                } else {
                    callback.onFailure(-1, permalink + "_basic_failure");
                }
            }
        }.start();
    }
}

这种方法很有效,但我不确定我是否以正确的方式做事。请看一下第 1-5 行:

  1. 从列表创建 Observable 时,我应该迭代吗 我自己在清单上,还是有其他更好的方法?为了 例如,有 Observable.from(Iterable),但我不认为我 可以用吗?
  2. 我在发送 REST 请求之前检查 isUnsubscribed,并且 也在两个结果处理程序中(成功/失败)。是这样吗 应该用吗?
  3. 我实现了一些逻辑来调用 onComplete 当所有请求都有 返回,无论成功或失败(见问题 5)。但是,就像他们一样 由 isUnsubscribed 调用保护,它们可能会发生 onComplete 永远不会被调用。流如何知道它 应该完成吗?当订阅者取消订阅并且我不必考虑它时,它是否会过早完成?有什么注意事项吗?例如,如果订阅者在发出内容结果的过程中取消订阅会发生什么?在我的测试中,从未发出过所有结果的列表,但我希望得到一个包含到目前为止所有结果的列表。
  4. 在 onFailure 方法中,我使用 onNext(null) 发出 null 来标记一个 失败。我这样做是因为最后我会有 2 个流的 zip, 并且只有在两者都压缩时才会发出自定义类的实例 值不为空。这是正确的做法吗?
  5. 正如我所提到的,我有一些自定义逻辑来检查流是否 完成的。我在这里所做的是计算 REST 结果以及何时作为 许多已处理,因为列表中有永久链接,已完成。 这是必要的吗?这是正确的方法吗?

【问题讨论】:

    标签: rx-java


    【解决方案1】:

    从列表创建 Observable 时,我应该自己迭代列表,还是有其他更好的方法?比如有Observable.from(Iterable),但我觉得用不上?

    您可以为每个链接构建一个可观察对象,而不是在 create 方法中进行迭代。

    public Observable<Content> getContent(permalink) {
    
        return Observable.create(subscriber -> {
                       ContentManager cm = new ContentManager();
                       if (!subscriber.isUnsubscribed()) { 
                                cm.getBasicContentByPermalink(permalink, new RestCallback() {
    
                                    @Override
                                    public void onSuccess(Content content) {
                                            subscriber.onNext(content);
                                            subscriber.onCompleted();
                                    }
    
                                    @Override
                                    public void onFailure(int code, String message) {
                                           subscriber.onError(OnErrorThrowable.addValueAsLastCause(new RuntimeException(message, permalink));
                                    }
    
                    }
        });
    }
    

    然后将其与flatMap 运算符一起使用

     Observable.from(permalinks)
               .flatMap(link -> getContent(link))
               .subscribe();
    

    我在发送 REST 请求之前检查 isUnsubscribed,也在两个结果处理程序(成功/失败)中检查。这就是它应该使用的方式吗?

    在发送请求之前检查 isUnsubscribed 是一种方法。不确定这个成功/失败回调(我认为它没用但有人可以说我错了)

    我实现了一些逻辑来在所有请求都返回时调用 onComplete,无论成功还是失败(参见问题 5)。但是,由于它们受到 isUnsubscribed 调用的保护,它们的 onComplete 可能永远不会被调用。

    如果您取消订阅直播,您将停止观看直播,因此您可能不会收到直播完成的通知。

    例如,如果订阅者在发出内容结果的过程中取消订阅会发生什么?

    您将能够使用退订前发出的内容。

    在 onFailure 方法中,我使用 onNext(null) 发出 null 来标记失败。我这样做是因为最后我将有 2 个流的 zip,并且只有当两个压缩值都不为空时才会发出自定义类的实例。这是正确的做法吗?

    不。而不是 null,发出错误(在订阅者上使用 onError)。 使用它,您不必在 zip lambda 中检查 null。只需压缩!

    正如我所提到的,我有一些自定义逻辑来检查流是否完成。我在这里所做的是计算 REST 结果,当处理的数量与列表中的永久链接一样多时,它就完成了。这是必要的吗?这是正确的方法吗?

    如果要获取永久链接的数量,当流完成时,可以使用count运算符。

    Observable.from(permalinks)
              .flatMap(link -> getContent(link))
              .count()
              .subscribe(System.out::println);
    

    如果你想发出链接数,在成功链接的末尾,你可以试试scan operator

     Observable.from(permalinks)
              .flatMap(link -> getContent(link))
              .map(content -> 1) // map content to 1 in order to count it.
              .scan((seed, acu) -> acu + 1) 
              .subscribe(System.out::println); // will print 2, 3...
    

    【讨论】:

    • 我自己尝试了 from/flatMap 但我犯了一个错误,你的代码帮助我正确地做到了,我更喜欢这种方式,不需要我的自定义逻辑来检测完成,谢谢!在行subscriber.onError(OnErrorThrowable.addValueAsLastCause(new RuntimeException("failed"), permalink));我得到:线程“Thread-7”rx.exceptions.OnErrorNotImplementedException 中的异常:失败,并且从未看到流中的任何排放。我不认为这是要走的路。我根本不会调用 onNext,然后调用 onCompleted,这会将永久链接转换为 0 内容。我什至不需要空值过滤器。
    • 您应该使用 onError 而不是 onCompleted 来通知错误。如果你想检测错误,你将实现 onError 回调:'Observable.error(new RuntimeException("oups").subscribe(System.out::println, (e) -> System.err.println( "ERROR ->" + e.getMessage());'。这样,您就可以管理流中的错误(使用 onErrorResumeNext 方法,...)
    • toList 只有在流完成时才会发出一个列表。由于它没有完成,它不会发出列表。如果你想要一个发射项目的部分列表,你必须将每个内容转换为一个项目的列表,然后使用扫描组合它。每次远程调用成功时,它都会发出一个列表。
    • 只处理错误然后: Observable.from(permalinks) .flatMap(link -> getContent(link).doOnError(e -> System.out.println("Get Exception : " + e) .onErrorResumeNext(Observable.empty())。这个例子浅显了网络调用主流的错误。因此流将继续并忽略错误。
    • getContent() 流将被停止。不是“主流”。 (请注意 onerrorresumenext 在 getContent() 流上!)。这是要做的事情。否则:如何知道是否遇到错误或流是否完成?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-11-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-04-28
    • 2011-03-17
    相关资源
    最近更新 更多