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

如何实现可周期性轮询至匹配断言的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:26:13