RxJava:实现发射+处理全周期Observable的方案咨询
解决方案
你之前的方案核心问题是takeUntil会在simulator.emitting发出false的瞬间直接截断序列,完全无视还有多少正在异步处理的数据项,导致未完成的处理结果丢失,也无法发出结束信号。要实现需求,需要完成三个核心动作:
- 先发出「开始发射并处理」的标识信号
- 接收并处理所有已发射的
items,直到emitting变为false后停止接收新项,但必须等已接收的项全部处理完成 - 所有项处理完毕后,发出「结束发射并处理」的标识信号
具体实现(以RxJava为例)
首先定义一个密封类来区分不同类型的事件:
sealed class ProcessingEvent { object StartProcessing : ProcessingEvent() data class ProcessedData<T>(val data: T) : ProcessingEvent() object EndProcessing : ProcessingEvent() }
然后构建目标Observable:
// 1. 生成开始信号:监听第一次发射开始的指令 val startSignal = simulator.emitting .filter { it } .take(1) .map { ProcessingEvent.StartProcessing } // 2. 处理数据项:接收items直到发射停止,异步处理每个项并发射结果 val processedItems = simulator.items .takeUntil(simulator.emitting.filter { !it }) // 停止接收新的items .concatMap { item -> // 指定处理线程,执行你的处理逻辑 Observable.just(item) .observeOn(Schedulers.io()) .map { process(it) } .map { ProcessingEvent.ProcessedData(it) } } // 3. 生成结束信号:等待发射停止,且所有已接收的项处理完成后再发射 val endSignal = simulator.emitting .filter { !it } .take(1) // 延迟发射结束信号,直到所有已接收的items处理完毕 .delaySubscription( simulator.items .takeUntil(simulator.emitting.filter { !it }) .concatMap { item -> Observable.just(item) .observeOn(Schedulers.io()) .map { process(it) } } .ignoreElements() // 忽略处理结果,仅等待处理完成 ) .map { ProcessingEvent.EndProcessing } // 拼接三个部分,得到最终的Observable val processingObservable = Observable.concat(startSignal, processedItems, endSignal)
关键逻辑说明
concatMap保证了数据项的处理顺序,且会等待前一个项处理完成后再处理下一个,避免并发处理的混乱(如果不需要严格顺序,可替换为flatMap,但需注意线程安全)delaySubscription确保结束信号只会在所有已接收的数据项处理完成后才发射takeUntil(simulator.emitting.filter { !it })确保不再接收新的数据项,但已经接收的会被完整处理
最终的processingObservable会严格按照你想要的序列发射事件:StartProcessing → 每个处理后的ProcessedData → EndProcessing
内容的提问来源于stack exchange,提问作者mediumkuriboh
相关产品推荐
相关产品推荐

