【问题标题】:RxJava subscription to side effectRxJava 订阅副作用
【发布时间】:2017-04-27 02:19:25
【问题描述】:

我有一个问题,就是简单的业务逻辑流程:

查看多个部门的员工,部门与员工的关系是否在缓存中,先查看缓存中是否存在关系,如果存在,查看是否存在 员工属于它,如果不在缓存中,则从数据库中获取,并检查与员工的关系,然后将部门信息保存到缓存中。

这是代码:

public Observable  isEmployeeInDepartment(List<Long> departmentIds, long employeeId){

     //this observable will resolve twice, and cause unnecessary cache access
     Observable  departmentInfoExsitInCache= checkDepartmentInfoFromCache(...).share();   

     Observable  departInfoNotInCache = departmentInfoExsitInCache.filter(...);

     //this observable will resolve twice, and cause unnecessary database access
     Observable  departmentInfoFromDb=departInfoNotInCache.flatMap(departmentIds->checkFromDb()).share(); 

     Observable<Long> saveResult=departmentInfoFromDb.flatMap(departmentInfo->saveToCache());

     Observable<Long> departInfoInCache = departmentInfoExsitInCache.filter(...);

     return departInfoInCache.check(userId).merge( departmentInfoFromDb.check(userId)).doOnCompleted(saveResult.subscribe());
}

问题是departmentInfoExsitInCache 和saveResult 会被客户端方法subscribe 解决两次。

我发现一旦删除保存订阅代码.doOnCompleted(saveResult.subscribe()),它就会变得正常并且只解决一次。这段代码有什么问题吗?

【问题讨论】:

    标签: rx-java reactive-programming side-effects


    【解决方案1】:

    您在此处滥用共享。
    问题是share() 在这种情况下不会帮助你。 share() 实际上是 publish().autoConnect() 用于保持对流的单一订阅,因此再次调用 subscribe 不会再次调用订阅逻辑,而只会将您连接到现有流。
    但是,在所有订阅者取消订阅共享流后,Observable 将取消订阅,这意味着当您再次调用subscribe() 时,您将调用订阅逻辑并再次调用 DB/Cache。

    所以,您正在做的是在取消订阅后再次订阅共享运算符。 (在doOnCompleted() 中),这将导致departmentInfoFromDbdepartmentInfoExsitInCache 再次订阅并转到数据库/缓存。

    考虑使用cache()/reply() 运算符在订阅之间持久化从数据库/缓存中获取的值。

    【讨论】:

    • 事实上,我不知道有什么方法可以让saveResult在整个方法返回的observable得到下标时被解析,除了把它放在“main”返回的observable的doOnCompleted()方法中让它也下标
    • 如果想在最后一行合并Obesrvable后订阅saveResult,可以使用concat()操作符。但无论如何它不会解决这个问题
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-03-19
    • 1970-01-01
    • 1970-01-01
    • 2021-03-25
    • 1970-01-01
    相关资源
    最近更新 更多