【问题标题】:RxJava Caching Results of Network Call IDs from a Stream (Redis or similar caching solution)来自流的网络调用 ID 的 RxJava 缓存结果(Redis 或类似的缓存解决方案)
【发布时间】:2015-12-09 18:26:06
【问题描述】:

我需要一些关于 RxJava 的帮助。我有一个昂贵的网络调用,它返回一个 Observable(来自 elasticsearch 的广告流)。我想将每个发出的项目(广告)的 ID 属性缓存 10 分钟(在 Redis 中),以便接下来 10 分钟内的后续调用使用缓存中的 ID 从 Elasticsearch 中获取广告。

我有一些代码 - 这在某种程度上有助于实现预期的结果,(感谢以下博客 .. http://blog.danlew.net/2015/06/22/loading-data-from-multiple-sources-with-rxjava/)

它可以缓存流中每个发出的项目,我需要的是它将流中这些项目的所有 ID 缓存为 1 个缓存条目

到目前为止的代码在这里 https://github.com/tonymurphy/rxjava 供任何感兴趣的人使用,下面是 sn-ps

@Component
public class CachingObservable {

    private final Logger logger = LoggerFactory.getLogger(CachingObservable.class);

    @Autowired
    private AdvertService advertService;

    // Each "network" response is different
    private Long requestNumber = 0L;

    public Observable<Advert> getAdverts(final String location) {

        Observable<Advert> memory = memory(location);
        Observable<Advert> network = network(location);

        Observable<Advert> networkWithSave = network.doOnNext(new Action1<Advert>() {
            @Override
            public void call(Advert advert) {
                List<Long> ids = new ArrayList<Long>();
                ids.add(advert.getId());
                advertService.cache(location, ids);
            }
        });

        // Retrieve the first source with data - concat checks in order
        Observable<Advert> source = Observable.concat(memory,
                networkWithSave)
                .first();

        return source;
    }

据我了解,concat 方法对我的用例并没有真正的用处。我需要知道网络可观察对象是否/何时完成,我需要获取返回的广告 ID 列表,并且需要将它们存储在缓存中。我可以订阅网络 observable - 但我希望它是懒惰的 - 只有在缓存中找不到数据时才调用。所以以下更新的代码不起作用..任何想法表示赞赏

public Observable<Advert> getAdverts(final String location) {

    Observable<Advert> memory = memory(location);
    Observable<Advert> network = network(location);

    Observable<Advert> networkWithSave = network.doOnNext(new Action1<Advert>() {
        @Override
        public void call(Advert advert) {
            List<Long> ids = new ArrayList<Long>();
            ids.add(advert.getId());
            advertService.cache(location, ids);
        }
    });

    // Retrieve the first source with data - concat checks in order
    Observable<Advert> source = Observable.concat(memory,
            networkWithSave)
            .first();

    Observable<List<Advert>> listObservable = networkWithSave.toList();
    final Func1<List<Advert>, List<Long>> transformer = new Func1<List<Advert>, List<Long>>() {
        @Override
        public List<Long> call(List<Advert> adverts) {
            List<Long> ids = new ArrayList<Long>();
            for (Advert advert : adverts) {
                ids.add(advert.getId());
            }
            return ids;
        }
    };

    listObservable.map(transformer).subscribe(new Action1<List<Long>>() {
        @Override
        public void call(List<Long> ids) {
            logger.info("ids {}", ids);
        }
    });

    return source;
}

【问题讨论】:

    标签: caching rx-java observable


    【解决方案1】:

    我要做的是使用过滤器来确保缓存的旧内容不会被发出,以便 concat 跳转到网络调用:

    Subject<Pair<Long, List<Advert>>, Pair<Long, List<Advert>>> cache = 
        BehaviorSubject.create().toSerialized();
    static final long RETENTION_TIME = 10L * 60 * 1000;
    
    Observable<Advert> memory = cache.filter(v -> 
        v.first + RETENTION_TIME > System.currentTimeMillis()).flatMapIterable(v -> v);
    
    Observable<Advert> network = ...
    
    Observable<Advert> networkWithSave = network.toList().doOnNext(v -> 
        cache.onNext(Pair.of(System.currentTimeMillis(), v)).flatMapIterable(v -> v)
    );
    
    return memory.switchIfEmpty(network);
    

    【讨论】:

      【解决方案2】:

      好的,我想我有一个适合我的解决方案。我可能忽略了一些东西,但这应该很容易吧?

      public Observable<Advert> getAdverts(final String location) {
      
          Observable<Advert> memory = memory(location);
          final Observable<Advert> network = network(location);
      
          final Func1<List<Advert>, List<Long>> advertToIdTransformer = convertAdvertsToIds();
      
          memory.isEmpty().subscribe(new Action1<Boolean>() {
              @Override
              public void call(Boolean aBoolean) {
                  if (aBoolean.equals(Boolean.TRUE)) {
                      Observable<List<Long>> listObservable = network.toList().map(advertToIdTransformer);
                      listObservable.subscribe(new Action1<List<Long>>() {
                          @Override
                          public void call(List<Long> ids) {
                              logger.info("Caching ids {}", ids);
                              advertService.cache(location, ids);
                          }
                      });
                  }
              }
          });
      
          // Retrieve the first source with data - concat checks in order
          Observable<Advert> source = Observable.concat(memory,
                  network)
                  .first();
      
      
          return source;
      }
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2011-03-01
        • 1970-01-01
        • 2010-11-05
        • 2013-01-12
        • 2017-09-26
        • 1970-01-01
        • 2017-08-17
        • 1970-01-01
        相关资源
        最近更新 更多