RxJava/RxKotlin实现列表扫描:找到全部目标即停止或超时报错
解决RxJava中找到所有目标后提前终止Observable的问题
你说得对,只添加另一个Predicate确实没法实现“找到全部目标就停止”的逻辑——因为单一的过滤操作只能判断当前元素是否匹配,没法跟踪已经收集到的目标集合,自然不知道什么时候已经找全了所有目标。我们需要结合RxJava的状态跟踪操作符来实现这个需求。
核心思路
- 把目标列表转换成
Set,提高“是否包含目标”和“是否找全”的判断效率 - 跟踪已找到的目标集合,每次收到匹配的元素就更新这个集合
- 当已找到的集合包含所有目标时,立即终止Observable
- 保留原有的超时逻辑,确保在超时前没找全目标时抛出错误
完整实现代码
import io.reactivex.rxjava3.core.Observable import java.util.concurrent.TimeUnit private const val TAG = "ScanManager" fun scanForValues(targetList: List<String>): Observable<FoundNumber> { val targetSet = targetList.toSet() val scanResult = scanForTargets() return scanResult // 打印所有扫描到的潜在目标 .doOnNext { Log.d(TAG, "Found potential target: ${it.targetId}") } // 过滤出非null且在目标列表中的项 .filter { it.targetId != null && it.targetId in targetSet } // 用scan跟踪已找到的目标集合,同时保留当前元素 .scan(Pair(mutableSetOf<String>(), null as FoundNumber?)) { (foundSet, _), current -> foundSet.add(current.targetId!!) Pair(foundSet, current) } // 当已找到的集合包含所有目标时,终止流 .takeUntil { (foundSet, _) -> foundSet.containsAll(targetSet) } // 提取当前找到的元素(忽略累积的集合) .map { (_, current) -> current!! } // 超时逻辑:5秒内没找全则抛出错误 .timeout(5, TimeUnit.SECONDS) .doOnError { Log.w(TAG, "Failed to scan: ${it.message}") } } // 你原有的scanForTargets及相关类 var emittedList: List<String?> = listOf(null, "0", "1", "2", "3") fun scanForTargets(): Observable<FoundNumber> = Observable .intervalRange(0, emittedList.size.toLong(), 0, 1, TimeUnit.SECONDS) .map { index -> FoundNumber(emittedList[index.toInt()]) } data class FoundNumber(val targetId: String?)
关键细节解释
scan操作符:用来维护一个可变的Set,每次收到匹配的元素就把它加入集合,同时返回“集合+当前元素”的配对。这样我们既能跟踪已找到的目标,又能保留需要发射的原始元素。takeUntil操作符:监听scan返回的配对,判断累积的集合是否包含所有目标。一旦满足条件,就立即终止整个Observable,实现“找全即停”的逻辑。Set的使用:目标列表转Set是为了快速判断元素是否在目标中;累积集合用Set是为了自动去重(避免重复收集同一个目标)。- 空值处理:添加了
it.targetId != null的判断,避免FoundNumber的null值干扰后续逻辑。
测试验证
当你调用scanForValues(listOf("1", "2"))时:
- 扫描流会依次发射
null、"0"、"1"、"2"、"3"对应的FoundNumber - 过滤后只保留"1"和"2"对应的元素
- 当发射完"2"时,累积的
Set变为{"1","2"},满足containsAll(targetSet)的条件,takeUntil触发,流立即终止 - 最终Observable会发射"1"和"2",不会等待5秒超时
内容的提问来源于stack exchange,提问作者user443654
相关产品推荐
相关产品推荐

