代码之家  ›  专栏  ›  技术社区  ›  AdrianS

将多个异步Observable<List>减少为一个Observatable<List〕

  •  1
  • AdrianS  · 技术社区  · 10 年前

    所以我想做的是:
    对于每个我想称之为 Service.getList() 10次,然后将10个列表合并为一个。

    现在我尝试了这两种方法,它们都在UT中工作,但在实际应用程序中失败(我猜我没有正确地执行reduce操作,因为在实际应用中存在异步http调用)。 对于这两种情况,我都看不到日志放在后面 reduce() 也不 onNext() 也没有 onError() ,所以我猜 reduce() 操作未完成。

    尝试#1:

    public Observable<List<Event>> getEventsForLocation(Location location) {
    List<Observable<List<Event>>> obs = new ArrayList<>();
    for (Venue v : location.getVenues()) {
         obs.add(getEventsForVenue(v)); //does one http call, returns Observable<List<Event>> 
    }
    return Observable.concat(Observable.from(obs))
                .reduce((List<Event>) new ArrayList<Event>(), (events, events2) -> {
                    events.addAll(events2);
                    return events;
                })
                .doOnNext(events -> Log.d("reduce ", events.toString()))
                .doOnError(throwable -> Log.e("reduce error", throwable.toString()));}
    

    尝试#2:

    public Observable<List<Event>> getEventsForLocation(Location location) {
    return Observable
            .from(location.getVenues())
            .flatMap(venue -> getEventsForVenue(venue)) //does one http call, returns Observable<List<Event>> 
            .reduce((List<Event>) new ArrayList<Event>(), (events, events2) -> {
                events.addAll(events2);
                return events;
            })
           .doOnNext(events -> Log.d("service", "total events " + events.toString()))
           .doOnError(t -> Log.e("service", "total events error2 " + t.toString()));}
    

    两种引道通过的UT:

        @Test
        public void getEventsForLocation() {
            Location loc = new Location("test", newArrayList(new Venue("v1", "url1"),new Venue("v2", "url2")));
    
            when(httpGateway.downloadWebPage(Mockito.anyString())).thenReturn(
                    Observable.just(readResource("eventsForVenue1.html")),
                    Observable.just(readResource("eventsForVenue2.html"))
            );
    
            TestSubscriber<List<Event>> probe = new TestSubscriber<>();
            service.getEventsForLocation(loc).subscribe(probe);
    
            probe.assertNoErrors();
    
            //assert the next event containts contents of all lists
            List<Event> events = probe.getOnNextEvents().get(0);
    
            //first list
            Assert.assertEquals("Unexpected title", "event1", events.get(0).getName());
            Assert.assertEquals("Unexpected artist", "artist1", events.get(0).getArtist());
            //second list
            Assert.assertEquals("Unexpected title", "event2", events.get(1).getName());
            Assert.assertEquals("Unexpected artist", "artist2", events.get(1).getArtist());
        }
    

    更新

    下面是更完整的代码和调度程序。

    Observable
                    .just(loc)
                    .subscribeOn(AndroidSchedulers.mainThread())
                    .observeOn(Schedulers.io())
                    .flatMap(location -> service.getEventsForLocation(location))
                    .observeOn(AndroidSchedulers.mainThread())
                    .subscribe(getObserver();
    
    3 回复  |  直到 7 年前
        1
  •  0
  •   akarnokd    10 年前

    如果您不关心列表中的最终订单,可以使用 from + flatMap + flatMapIterable + toList :

    Observable.from(location.getVenues())
    .flatMap(venue -> getEventsForVenue(venue))
    .flatMapIterable(list -> list)
    .toList();
    

    如果顺序很重要,并且您希望“并行”执行getEventsForVenue,则可以将flatMap替换为concatMapEager:

    Observable.from(location.getVenues())
    .concatMapEager(venue -> getEventsForVenue(venue))
    .concatMapIterable(list -> list)
    .toList();
    
        2
  •  0
  •   paul    10 年前

    你可以使用Collect,它会成功的。但在这种情况下,我合并了几个项目,通常我更喜欢使用扫描。对于每个新项目,它将为您提供最后一个处理的项目。因此,您可以将每个项目附加到前面的项目。

    看看大理石图,以防箱子适合您的箱子 http://reactivex.io/documentation/operators/scan.html

    尝试从这个简单的例子开始,然后尝试将其应用到代码中

    /**
     * apply this function for every item against the previous emitted item from the source.
     *  Emitted:
                0
                1
                3
                6
                10
                15
     */
    @Test
    public void scanObservable() {
        Integer[] numbers = {0, 1, 2, 3, 4, 5};
    
        Observable.from(numbers)
                  .scan((lastItemEmitted, newItem) -> (lastItemEmitted + newItem))
                  .subscribe(System.out::println);
    }
    
        3
  •  0
  •   AdrianS    10 年前

    我终于找到了一个解决方案,使用zip(),但我不喜欢它。
    这个问题应该能够通过flatMap/reduce的组合来解决

    public Observable<List<Event>> getEventsForLocation(Location location) {
            List<Observable<List<Event>>> venues = new ArrayList<>();
            for (Venue v : location.getVenues()) {
                venues.add(Observable.just(v).flatMap(venue -> getEventsForVenue(venue)));
            }
            return Observable.zip(venues, new FuncN<List<Event>>() {
                @Override
                public List<Event> call(Object... args) {
                    List<Event> allEvents = new ArrayList<Event>();
                    for (Object o : args) {
                        List<Event> le = (List<Event>) o;
                        allEvents.addAll(le);
                    }
                    return allEvents;
                }
            });