RxJava下满足动态停止条件的Reduce Single实现及优化问题
最终优化实现
fun readAllLogs(): Single<LogEntry> { return Single.just(0) // 包装请求结果,同时携带累计偏移量 .concatMapSingle { position -> readBytes(position).map { bytes -> position to bytes } } // 基于上一次结果递归发起下一次请求,流内部维护偏移量状态 .expand { (prevTotal, bytes) -> if (bytes.isEmpty()) { // 读到空数组直接停止后续请求 Single.empty() } else { val newPosition = prevTotal + bytes.size readBytes(newPosition).map { newBytes -> newPosition to newBytes } } } // 提取返回的字节数组,过滤终止用的空数组 .map { it.second } .takeWhile { it.isNotEmpty() } // 累加结果返回Single .reduceWith({ LogEntry() }) { logEntry, bytes -> logEntry += bytes logEntry } }
问题1:如何优雅地终止流,无需在流外部维护独立的状态变量?
用RxJava标准操作符expand即可实现完全无外部状态的流式调用:
- 所有偏移量状态都通过流事件传递,保存在流内部,不需要外部定义原子类或者其他变量存储状态,也不存在线程安全问题
- 终止逻辑直接通过
takeWhile判断返回的字节数组是否为空即可触发,不需要额外对比新旧偏移量是否相等,逻辑更直观
问题2:两次在Single和Observable之间转换,有没有更高效的优化方案?
上面的实现已经完全去掉了不必要的类型转换:
- 全程除了最终
reduce返回Single之外,没有任何多余的Observable和Single互转操作,减少了不必要的事件包装开销 - 操作符链逻辑线性顺承,可读性和可维护性都更高,性能对比原有实现也有提升
内容的提问来源于stack exchange,提问作者Entreco
相关产品推荐
相关产品推荐

