Scala Monix问题:Observable.foldLeft结果无法传入flatMap
问题分析与解决
核心问题
你用的foldLeft是Rx Observable的终止操作符——它只有在源Observable完全结束(发送onComplete信号)时,才会把最终的累加结果发射给后续操作。但你的源Observable.repeatEval(Random.nextInt(10))是一个无限流,永远不会主动结束,所以foldLeft永远不会向flatMap传递任何数据,这就是为什么flatMap里的代码完全没执行。
另外你原代码里的逻辑还有个小问题:当累加值超过15后,你只是打印但返回原累加值,导致后续所有新的随机值都不会再被累加,累加器会一直停在≥15的状态,完全失去了继续合并的能力。
修正方案
改用scan操作符——它会每次发射累加的中间结果,刚好符合你需要实时监控累加值、触发阈值处理的需求。同时调整逻辑,超过阈值后重置累加器,继续收集新的随机值,最后用filter只把符合条件的结果传递给后续处理器。
修正后的代码:
import rx.lang.scala.Observable import scala.util.Random // 生成0-9的随机值无限流 val randomStream: Observable[Int] = Observable.repeatEval(Random.nextInt(10)) // 合并值+阈值处理 val processedStream = randomStream .scan(0) { (acc, current) => val newAcc = acc + current if (newAcc >= 15) { println("handled " + newAcc) 0 // 超过阈值后重置累加器,继续收集下一批值 } else { newAcc } } // 只保留超过阈值的结果,传递给后续flatMap .filter(_ >= 15) .flatMap { handledValue => val mappedValue = handledValue + 1 println("mapped " + mappedValue) Observable(mappedValue) } // 订阅触发流执行(冷流必须订阅才会开始工作) processedStream.subscribe()
代码说明
scan的作用:每次接收新的随机值后,计算新的累加值并立即发射这个中间结果,让我们能实时判断是否触发阈值。- 累加器重置:当累加值≥15时,打印后重置为0,这样后续的随机值会重新开始累加,保证流程持续运行。
filter的作用:过滤掉未超过阈值的中间累加值,只把真正处理过的阈值结果传给flatMap,避免无效数据干扰。- 订阅触发:Rx的Observable是冷流,必须调用
subscribe()才会启动整个流的执行。
内容的提问来源于stack exchange,提问作者Jelly
相关产品推荐
相关产品推荐

