You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.30 04:27:00