如何实现可周期性轮询至匹配断言的RxObservable?
解决方案:实现带断言终止的周期性轮询Observable
我来帮你完善这个Watcher类,满足你提出的所有轮询需求,直接上修改后的完整代码和详细说明:
import io.reactivex.Observable import io.reactivex.disposables.Disposable import java.util.concurrent.ConcurrentHashMap import java.util.concurrent.TimeUnit class Watcher<T : Any>( /** * Emits the data associated with the provided id */ private val callable: (id: String) -> T, /** * Checks if the provided value marks the observable as complete */ private val predicate: (id: String, value: T) -> Boolean, // 添加轮询间隔参数(单位:毫秒) private val pollInterval: Long ) { private val watchPool: MutableMap<String, Observable<T>> = ConcurrentHashMap() fun watch(id: String): Observable<T> { // 复用已存在的Observable,避免同一id重复轮询 val existing = watchPool[id] if (existing != null) return existing val initialValue = callable(id) // 如果初始值就满足终止断言,直接返回单次发射的Observable,不存入池 if (predicate(id, initialValue)) { return Observable.just(initialValue) } // 创建具备完整轮询逻辑的Observable val observable = Observable.defer { // 每次轮询都调用callable获取最新值 Observable.just(callable(id)) } .repeatWhen { completions -> // 每次轮询完成后,延迟指定间隔再触发下一次轮询 completions.flatMap { Observable.timer(pollInterval, TimeUnit.MILLISECONDS) } } .takeUntil { value -> // 检查当前值是否满足终止条件,满足则立即结束序列 predicate(id, value) } .doOnDispose { // 当所有订阅者取消订阅时,从池中移除该id的Observable watchPool.remove(id) } .doOnComplete { // 序列因断言满足正常完成时,从池中移除该id的Observable watchPool.remove(id) } .distinctUntilChanged() // 仅发射与前一次不同的值,避免重复通知 .share() // 共享订阅:同一id仅维护一个轮询任务,无订阅者时自动停止轮询 // 将创建好的Observable存入池,供后续复用 watchPool[id] = observable return observable } }
核心需求对应实现说明
这里逐个拆解你提出的需求,说明代码是如何满足的:
有订阅者才开始轮询:
结合Observable.defer和share()操作符实现。defer会延迟创建发射源,直到有订阅者订阅;share()会管理订阅计数,只有当第一个订阅者出现时,才真正启动上游的轮询序列,避免无意义的后台轮询。每x秒执行一次轮询:
在repeatWhen中使用Observable.timer(pollInterval, TimeUnit.MILLISECONDS),每次轮询完成(即调用callable并发射值后),等待指定间隔再触发下一次轮询,精准控制轮询频率。断言返回true时标记为完成状态:
使用takeUntil操作符,每次发射值后都会检查predicate(id, value),一旦断言返回true,takeUntil会立即终止整个序列,触发onComplete事件,同时触发doOnComplete从池中移除该id的Observable。订阅者数量从大于0变为0时自动完成:
share()操作符会跟踪订阅者数量,当最后一个订阅者取消订阅时,上游的轮询序列会被自动dispose,此时doOnDispose会触发,将该id对应的Observable从watchPool中移除,轮询任务随之停止。
示例场景使用(Stage枚举轮询)
针对你提到的订单状态轮询场景,这里给出完整的使用示例:
enum class Stage { CREATED, PROCESSING, DELIVERING, FINISHED } // 模拟根据订单ID获取当前状态的业务逻辑 fun getOrderStage(orderId: String): Stage { // 这里替换为实际的查询逻辑(比如数据库、接口调用) val currentTime = System.currentTimeMillis() return when { currentTime % 10000 < 3000 -> Stage.CREATED currentTime % 10000 < 6000 -> Stage.PROCESSING currentTime % 10000 < 9000 -> Stage.DELIVERING else -> Stage.FINISHED } } fun main() { // 创建Watcher实例:轮询间隔2秒,终止条件为状态变为FINISHED val orderWatcher = Watcher( callable = ::getOrderStage, predicate = { _, stage -> stage == Stage.FINISHED }, pollInterval = 2000 ) // 订阅订单ID为"order_123"的状态轮询 val disposable: Disposable = orderWatcher.watch("order_123") .subscribe( { stage -> println("订单状态更新:$stage") }, { error -> println("轮询出错:${error.message}") }, { println("订单已完成,轮询终止") } ) // 如果需要主动停止轮询,调用dispose即可 // disposable.dispose() }
内容的提问来源于stack exchange,提问作者Allan W
相关产品推荐
相关产品推荐

