【问题标题】:Delay except the First time Rxjava Android延迟除了第一次 Rxjava Android
【发布时间】:2016-11-14 22:55:20
【问题描述】:

我正在进行异步调用,10 秒后 1 分钟,这意味着将进行大约 6 次调用,但问题在于我希望它在特定的 condition 上应用 delay

Observable
.just(listOfSomethings_Locally)
.take(1, TimeUnit.MINUTES)
.serialize()
.delaySubscription( // this is confusing part 
() ->
    Observable.just(listOfItems_Network).take(10,TimeUnit.SECONDS)
) 

我想要的是将网络呼叫延迟 10 秒,第一次呼叫除外,并在 10 秒后取消网络呼叫,所以我应该在 1 分钟内有准确的 6 个呼叫。

编辑

由于场景混乱,这里重新定义场景:

我有大量本地驱动程序,我想发送 每 10 秒后向他们每个人发出请求,然后听另一个 订户检查司机是否在 10 秒内取消了它, 这个过程将持续大约 1 分钟,如果一个司机取消我应该 立即向下一个发送请求

到目前为止编写的代码:

Observable.from(driversGot)
                .take(1,TimeUnit.MINUTES)
                .serialize()
                .map(this::requestRydeObservable) // requesting for single driver from driversGot (it's a network call)
                .flatMap(dif ->
                        Observable.amb(
                                kh.getFCM().driverCanceledRyde(), // listen for if driver cancel request returns integer
                                kh.getFCM().userRydeAccepted()) // listen for driver accept returns RydeAccepted object
                                .map(o -> {
                                    if (o instanceof Integer) {
                                        return new RydeAccepted();
                                    } else if (o instanceof RydeAccepted) {
                                        return (RydeAccepted) o;
                                    }
                                    return null;
                                }).delaySubscription(10,TimeUnit.SECONDS)
                )
                .subscribeOn(Schedulers.io())
                .observeOn(AndroidSchedulers.mainThread())
                .subscribe(fua -> {
                    if (fua == null) {
                        UiHelpers.showToast(context, "Invalid Firebase response");
                    } else if (!fua.getStatus()) { // ryde is canceled because object is empty
                        UiHelpers.showToast(context, "User canceled ryde");
                    } else { // ryde is accepted
                        UiHelpers.showToast(context, "User accepted ryde");
                    }
                }, t -> {
                    t.printStackTrace();
                    UiHelpers.showToast(context,"Error sending driver requests");
                }, UiHelpers::stopLoading);

【问题讨论】:

  • 请提供有关您的可观察对象和用例的更多详细信息。据我了解您的更新:您有一个项目列表,它将被转换为一个 Observable。您想一次处理一个元素,该元素已推送给您。必须在每 10 秒内为您安排一次值。示例:Sec 0:值 1,Sec 10:值 2。每个发出的值都将通过网络调用进行处理。如果该项目在 10 秒内被取消,你检查另一个 obs。如果不是,一个元素的过程将是 60 秒。如果在 10 秒内取消,你会开始排队的下一个吗?
  • 您完全了解流程,除了 60 秒是列表中所有元素的总体限制时间
  • 那么,您的列表中有 6 个元素,或者如果您每 10 秒调度一个元素,您希望如何在 60 秒内完成?
  • 每个元素都有 10 秒的报价,但可以在其中取消,因此下一个元素应立即安排
  • 好的,知道了。能否请您提供一些关于 driverCanceledRyde、userRydeAccepted、requestRydeObservable、fua 的实现细节(返回类型)

标签: android multithreading rx-java rx-android rx-java2


【解决方案1】:

对您的代码的反馈

您不需要takeserialize,因为just 会立即+ 连续发出内容。

delaySubscription 似乎是一个奇怪的选择,因为在传递的 observable 生成事件之后,不会延迟进一步的事件(这与您的情况相矛盾)

选项 #1,仅 rx

使用delay + 计算其余事件的各个延迟(因此第一个延迟 0 秒,第二个延迟 1 秒,第三个延迟 3,...)

            AtomicLong counter = new AtomicLong(0);
    System.out.println(new Date());
    Observable.just("1", "2", "3", "4", "5", "6")
        .delay(item -> Observable.just(item).delay(counter.getAndIncrement(), TimeUnit.SECONDS))
        .subscribe(new Consumer<String>() {
            public void accept(String result) throws Exception {
                System.out.println(result + " " + new Date());
            }
        });        
        System.in.read();

选项 #2:速率限制

您的用例似乎适合速率限制,因此我们可以使用来自 guava 的 RateLimiter:

            RateLimiter limiter = RateLimiter.create(1);
    System.out.println(new Date());
    Observable.just("1", "2", "3", "4", "5", "6")
        .map(r -> {
            limiter.acquire();
            return r;
        })
        .subscribe(new Consumer<String>() {
            public void accept(String result) throws Exception {
                System.out.println(result + " " + new Date());
            }
        });        
        System.in.read();

两者的工作方式相似:

Tue Nov 15 11:14:34 EET 2016
1 Tue Nov 15 11:14:34 EET 2016
2 Tue Nov 15 11:14:35 EET 2016
3 Tue Nov 15 11:14:36 EET 2016
4 Tue Nov 15 11:14:37 EET 2016
5 Tue Nov 15 11:14:38 EET 2016
6 Tue Nov 15 11:14:39 EET 2016

如果您提出要求,限速器会更好地工作,例如处理 5 秒,那么它将允许下一个请求更快地进行以补偿延迟并达到 1req/s 的目标 10 秒。

【讨论】:

  • 在您的选项#1 中,如何将请求限制为 1 分钟时间,因为每个请求都在 (order number)-1 秒后执行,但它不符合我的问题,我所拥有的是大列表本地司机,我想每 10 秒后向他们每个人发送请求,并听取另一个订阅者检查司机是否在 10 秒内没有取消它,这个过程将持续大约 1 分钟,如果一个司机取消我应该立即发送请求到下一个
【解决方案2】:

您好,您可以延迟重试,有一种方法可以做到这一点here

【讨论】:

  • 这里有点棘手,提供的链接仅解释delay的行为
【解决方案3】:

我想更新@Ivans 的帖子,因为它缺少错误处理并且正在使用副作用。

这篇文章将只使用 RxJava-operators。 observable 将每 10 秒提供一个值。请求将在 10 秒后超时。如果超时,将返回一个备用值。

第一个测试方法将在 60 秒内收到 10 个值。如果最后一个请求早于 10 秒完成,则 observable 可能会在 60 秒之前完成。

public class TimeOutTest {
    private static String DUMMY_VALUE = "ERROR";

    @Test
    public void handles_in_60_seconds() throws Exception {
        List<Integer> actions = Arrays.asList(1, 2, 3, 4, 5, 6, 7, 8, 9, 10);

        Observable<Long> timer = Observable.interval(0, 10, TimeUnit.SECONDS).take(6);

        Observable<Integer> vales = Observable.fromIterable(actions)
                .take(6);

        Observable<Integer> observable = Observable.zip(timer, vales, (time, result) -> {
            return result;
        });

        Observable<String> stringObservable = observable.flatMap(integer -> {
            return longNetworkLong(9_000)
                    .timeout(10, TimeUnit.SECONDS)
                    .onErrorReturnItem(DUMMY_VALUE);
        }).doOnNext(s -> System.out.println("VALUE"));

        stringObservable.test()
                .awaitDone(60, TimeUnit.SECONDS)
                .assertValueCount(6);
    }

    @Test
    public void last_two_values_timeOut() throws Exception {
        List<Integer> actions = Arrays.asList(1, 2, 3, 4, 5, 6, 7, 8, 9, 10);

        Observable<Long> timer = Observable.interval(0, 10, TimeUnit.SECONDS).take(6);

        Observable<Integer> vales = Observable.fromIterable(actions)
                .take(6);

        Observable<Integer> observable = Observable.zip(timer, vales, (time, result) -> {
            return result;
        });

        Observable<String> stringObservable = observable
                .map(integer -> integer * 2500)
                .flatMap(integer -> {
                    return longNetworkLong(integer)
                            .timeout(10, TimeUnit.SECONDS)
                            .doOnError(throwable -> System.out.print("Timeout hit?"))
                            .onErrorReturnItem(DUMMY_VALUE);
                })
                .doOnNext(s -> System.out.println("VALUE"))
                .filter(s -> !Objects.equals(s, DUMMY_VALUE));

        stringObservable.test()
                .awaitDone(60, TimeUnit.SECONDS)
                .assertValueCount(4);

    }

    private Observable<String> longNetworkLong(int delayTime) {
        return Observable.fromCallable(() -> {
            Thread.sleep(delayTime);
            return "result";
        });
    }
}

【讨论】:

  • 它接受一个 List 并将其转换为一个 Observable。还有另一个 Observable 每 10 秒发出一个值。我从两者中取 6 个值并将它们压缩在一起。因此,我将每 10 秒收到一个推送给我的值,从秒 0 开始。对于每个值,我调用一个长时间运行的方法:longNetworkLong。如果 longNetworkLong Observable 的结果需要超过 10 秒才能进入,我会取消请求并提供一个备用值。我会看看你的新要求。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-05-15
相关资源
最近更新 更多