【问题标题】:Aggregate resource requests & dispatch responses to each subscriber聚合资源请求并向每个订阅者发送响应
【发布时间】:2016-04-29 21:51:19
【问题描述】:

我对 RxJava 还很陌生,并且在一个对我来说似乎很常见的用例中苦苦挣扎:

从应用程序的不同部分收集多个请求,聚合它们,进行单个资源调用并将结果分派给每个订阅者。

我尝试了很多不同的方法,使用主题、可连接的可观察对象、延迟的可观察对象......到目前为止,没有一个能成功。

我对这种方法非常乐观,但事实证明它和其他方法一样失败了:

    //(...)
    static HashMap<String, String> requests = new HashMap<>();
    //(...)

    @Test
    public void myTest() throws InterruptedException {
        TestScheduler scheduler = new TestScheduler();
        Observable<String> interval = Observable.interval(10, TimeUnit.MILLISECONDS, scheduler)
                .doOnSubscribe(() -> System.out.println("new subscriber!"))
                .doOnUnsubscribe(() -> System.out.println("unsubscribed"))
                .filter(l -> !requests.isEmpty())
                .doOnNext(aLong -> System.out.println(requests.size() + " requests to send"))
                .flatMap(aLong -> {
                    System.out.println("requests " + requests);
                    return Observable.from(requests.keySet()).take(10).distinct().toList();
                })
                .doOnNext(strings -> System.out.println("calling aggregate for " + strings + " (from " + requests + ")"))
                .flatMap(Observable::from)
                .doOnNext(s -> {
                    System.out.println("----");
                    System.out.println("removing " + s);
                    requests.remove(s);
                })
                .doOnNext(s -> System.out.println("remaining " + requests));

        TestSubscriber<String> ts1 = new TestSubscriber<>();
        TestSubscriber<String> ts2 = new TestSubscriber<>();
        TestSubscriber<String> ts3 = new TestSubscriber<>();
        TestSubscriber<String> ts4 = new TestSubscriber<>();

        Observable<String> defer = buildObservable(interval, "1");
        defer.subscribe(ts1);
        Observable<String> defer2 = buildObservable(interval, "2");
        defer2.subscribe(ts2);
        Observable<String> defer3 = buildObservable(interval, "3");
        defer3.subscribe(ts3);
        scheduler.advanceTimeBy(200, TimeUnit.MILLISECONDS);
        Observable<String> defer4 = buildObservable(interval, "4");
        defer4.subscribe(ts4);

        scheduler.advanceTimeBy(100, TimeUnit.MILLISECONDS);
        ts1.awaitTerminalEvent(1, TimeUnit.SECONDS);
        ts2.awaitTerminalEvent(1, TimeUnit.SECONDS);
        ts3.awaitTerminalEvent(1, TimeUnit.SECONDS);
        ts4.awaitTerminalEvent(1, TimeUnit.SECONDS);

        ts1.assertValue("1");
        ts2.assertValue("2"); //fails (test stops here)
        ts3.assertValue("3"); //fails
        ts4.assertValue("4"); //fails


    }

    public Observable<String> buildObservable(Observable<String> interval, String key) {

        return  Observable.defer(() -> {
                            System.out.printf("creating observable for key " + key);
                            return Observable.create(subscriber -> {
                                requests.put(key, "xxx");
                                interval.doOnNext(s -> System.out.println("filtering : key/val  " + key + "/" + s))
                                        .filter(s1 -> s1.equals(key))
                                        .doOnError(subscriber::onError)
                                        .subscribe(s -> {
                                            System.out.println("intern " + s);
                                            subscriber.onNext(s);
                                            subscriber.onCompleted();
                                            subscriber.unsubscribe();
                                        });
                            });
                        }
                )
                ;
    }

输出:

creating observable for key 1new subscriber!
creating observable for key 2new subscriber!
creating observable for key 3new subscriber!
3 requests to send
requests {3=xxx, 2=xxx, 1=xxx}
calling aggregate for [3, 2, 1] (from {3=xxx, 2=xxx, 1=xxx})
----
removing 3
remaining {2=xxx, 1=xxx}
filtering : key/val  1/3
----
removing 2
remaining {1=xxx}
filtering : key/val  1/2
----
removing 1
remaining {}
filtering : key/val  1/1
intern 1
creating observable for key 4new subscriber!
1 requests to send
requests {4=xxx}
calling aggregate for [4] (from {4=xxx})
----
removing 4
remaining {}
filtering : key/val  1/4

在第二个断言处测试失败(ts2 没有收到“2”) 结果证明伪聚合按预期工作,但值没有分派给相应的订阅者(只有第一个订阅者接收它)

知道为什么吗?

另外,我觉得我在这里错过了明显的东西。如果您想出更好的方法,我非常愿意听到。

编辑:添加一些关于我想要实现的内容。

我有一个 REST API 通过多个端点(例如用户/{userid})公开数据。此 API 还可以聚合请求(例如,用户/用户 1 和用户/用户 2)并在一个 http 请求中获取相应的数据,而不是两个。

我的目标是能够在给定的时间范围内(比如 10 毫秒)自动聚合来自应用程序不同部分的请求,最大批量大小(比如 10),发出聚合 http 请求,然后分派结果给相应的订阅者。

类似这样的:

// NOTE: those calls can be fired from anywhere in the app, and randomly combined. The timing and order is completely unpredictable

//ts : 0ms
api.call(userProfileRequest1).subscribe(this::show); 
api.call(userProfileRequest2).subscribe(this::show);

//--> 10 毫秒后,应该用这 2 个调用触发一个 http 聚合请求,映射响应项并将它们发送给相应的订阅者(这将显示正确的用户配置文件)

//ts : 20ms
api.call(userProfileRequest3).subscribe(this::show); 
api.call(userProfileRequest4).subscribe(this::show);
api.call(userProfileRequest5).subscribe(this::show); 
api.call(userProfileRequest6).subscribe(this::show);
api.call(userProfileRequest7).subscribe(this::show); 
api.call(userProfileRequest8).subscribe(this::show);
api.call(userProfileRequest9).subscribe(this::show); 
api.call(userProfileRequest10).subscribe(this::show);
api.call(userProfileRequest11).subscribe(this::show); 
api.call(userProfileRequest12).subscribe(this::show);

//--> should fire a single http aggregate request RIGHT AWAY (we hit the max batch size) with the 10 items, map the response items & send them to the corresponding subscribers (that will show the right user profile)   

我编写的测试代码(仅包含字符串)并粘贴在这个问题的顶部,旨在作为我最终实现的概念证明。

【问题讨论】:

    标签: aggregate rx-java observable dispatch


    【解决方案1】:

    您的Observable 构造不正确

     public Observable<String> buildObservable(Observable<String> interval, String key) {
    
        return interval.doOnSubscribe(() -> System.out.printf("creating observable for key " + key))
                       .doOnSubscribe(() -> requests.put(key, "xxx"))
                       .doOnNext(s -> System.out.println("filtering : key/val  " + key + "/" + s))
                       .filter(s1 -> s1.equals(key));
            }
    

    当你在 subscribersubsribe 时:这是一个糟糕的设计。

    我不知道你想要实现什么,但我认为我的代码应该与你的非常接近。

    请注意,对于所有副作用,我使用doMethods(如doOnNextdoOnSubscribe)来表明我明确表示我想做一个副作用。

    我通过直接返回 interval 替换您的 defer 调用:因为您想在您的 defer 调用中的自定义 observable 构建中发出所有 interval 事件,返回 interval observable 更好。

    请注意,您过滤了interval Observable:

    Observable<String> interval = Observable.interval(10, TimeUnit.MILLISECONDS, scheduler)
                .filter(l -> !requests.isEmpty()).
                // ... 
    

    所以,一旦你将一些东西放入 requests 映射,interval 就会停止发射。

    我不明白您想通过请求映射实现什么,但请注意,您可能希望避免副作用,并且更新此映射显然是一种副作用。

    关于 cmets 的更新

    您可能希望使用buffer 运算符来聚合请求,然后以批量方式执行请求:

        PublishSubject<String> subject = PublishSubject.create();
    
    
        TestScheduler scheduler = new TestScheduler();
    
        Observable<Pair> broker = subject.buffer(100, TimeUnit.MILLISECONDS, 10, scheduler)
                                         .flatMapIterable(list -> list) // you can bulk calls here
                                         .flatMap(id -> Observable.fromCallable(() -> api.call(id)).map(response -> Pair.of(id, response)));
    
        TestSubscriber<Object> ts1 = new TestSubscriber<>();
        TestSubscriber<Object> ts2 = new TestSubscriber<>();
        TestSubscriber<Object> ts3 = new TestSubscriber<>();
        TestSubscriber<Object> ts4 = new TestSubscriber<>();
    
        broker.filter(pair -> pair.id.equals("1")).take(1).map(pair -> pair.response).subscribe(ts1);
        broker.filter(pair -> pair.id.equals("2")).take(1).map(pair -> pair.response).subscribe(ts2);
        broker.filter(pair -> pair.id.equals("3")).take(1).map(pair -> pair.response).subscribe(ts3);
        broker.filter(pair -> pair.id.equals("4")).take(1).map(pair -> pair.response).subscribe(ts4);
    
        subject.onNext("1");
        subject.onNext("2");
        subject.onNext("3");
    
        scheduler.advanceTimeBy(1, TimeUnit.SECONDS);
    
        ts1.assertValue("resp1");
        ts2.assertValue("resp2");
        ts3.assertValue("resp3");
        ts4.assertNotCompleted();
    
        subject.onNext("4");
        scheduler.advanceTimeBy(1, TimeUnit.SECONDS);
        ts4.assertValue("resp4");
        ts4.assertCompleted();
    

    如果你想执行网络请求折叠,你可能需要检查 Hystrix :https://github.com/Netflix/Hystrix

    【讨论】:

    • 感谢您的回复!我编辑了我的问题,以添加更多关于我想要实现的目标的上下文。我的代码中的defer 是为了懒惰地创建可观察对象,所以我们肯定会使用requests 的内容,因为它是在订阅时。我经常面临的问题(如果我使用您的代码,我想我会再次遇到)是我希望我的订阅者只获得正确的价值,然后完成(或取消订阅)。流程的聚合部分(此处为“调用聚合”)将返回一组项目(例如 1,2、3,4),然后我需要将其分派给订阅者(1 用于订阅者 1,2 用于订阅者 2)
    • 基本上,我希望订阅者只收到他们正在等待的值,然后取消订阅,这样它就不会收到其他任何东西。一种事件总线,每个订阅者只接受他明确订阅的一个值。
    • 它应该与 PublishSubject 和缓冲区运算符一起使用。 (我会更新我的答案)
    • 我更新了我的答案。请注意,Pair 对象只是保持请求与响应关联的一对。您的问题非常接近 Hystrix 可以实现的目标。
    • 看起来像一个魅力我实际上有一些与此非常相似的东西,除了我没有使用 take() 运算符。我真的不知道它会完成!非常感谢!
    猜你喜欢
    • 1970-01-01
    • 2015-06-30
    • 1970-01-01
    • 1970-01-01
    • 2010-12-21
    • 1970-01-01
    • 2019-06-16
    • 1970-01-01
    • 2021-08-21
    相关资源
    最近更新 更多