【发布时间】: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