RxJava中如何基于Cold Flowable实现双条件匹配、发射元素并持续搜索?
在RxJava中实现"A元素缓存+后续B元素配对处理"的方案
嘿,这个需求其实是RxJava里很常见的序列状态跟踪+配对处理场景,我来给你一步步拆解实现思路和代码。
核心思路
我们需要维护一个「是否已缓存符合条件A的元素」的状态,然后按顺序遍历流:
- 没缓存A时,找到第一个符合A的元素就缓存它;
- 已缓存A时,寻找符合B的元素,找到后就用这对元素生成新结果并发射;
- 持续这个逻辑,直到流结束。
这里的关键是用scan操作符来跟踪状态——它可以帮我们在遍历流的过程中,把上一步的状态传递到下一步,非常适合这种需要记忆前序信息的场景。
代码实现
1. 状态容器说明
我们用一个匿名内部类来封装状态(你也可以定义正式的POJO类):
cachedA:缓存的符合条件A的元素(用Optional处理空值);emissions:当前步骤要发射的处理结果(每次清空后重新添加)。
2. 完整代码示例
假设我们的流元素是字符串,条件A是前缀为"A",条件B是前缀为"B",处理函数是把A和B拼接成A-B的形式:
import io.reactivex.rxjava3.core.Flowable; import java.util.ArrayList; import java.util.List; import java.util.Optional; import java.util.function.BiFunction; import java.util.function.Predicate; public class RxPairingExample { public static void main(String[] args) { // 模拟Cold Flowable数据源 Flowable<String> source = Flowable.just("A1", "X", "B1", "B2", "A2", "Y", "B3"); // 定义条件和处理函数 Predicate<String> isConditionA = s -> s.startsWith("A"); Predicate<String> isConditionB = s -> s.startsWith("B"); BiFunction<String, String, String> combineElements = (a, b) -> a + "-" + b; source // 用scan维护状态,跟踪缓存的A元素和待发射结果 .scan(new Object() { Optional<String> cachedA = Optional.empty(); List<String> emissions = new ArrayList<>(); }, (accumulator, currentElement) -> { // 清空上一步的待发射列表,准备处理当前元素 accumulator.emissions.clear(); if (accumulator.cachedA.isEmpty()) { // 还没缓存A,检查当前元素是否符合条件A if (isConditionA.test(currentElement)) { accumulator.cachedA = Optional.of(currentElement); } } else { // 已经缓存了A,检查当前元素 if (isConditionB.test(currentElement)) { // 找到符合B的元素,生成结果并加入待发射列表 accumulator.emissions.add( combineElements.apply(accumulator.cachedA.get(), currentElement) ); // 👉 可选逻辑:如果需要每配对一次就重置缓存(重新找下一个A),就取消下面的注释 // accumulator.cachedA = Optional.empty(); } else if (isConditionA.test(currentElement)) { // 遇到新的符合A的元素,替换缓存的A accumulator.cachedA = Optional.of(currentElement); } } return accumulator; }) // 把每个步骤的待发射列表展开成流 .flatMapIterable(accumulator -> accumulator.emissions) // 订阅输出结果 .subscribe(System.out::println); } }
两种行为的结果对比
- 不重置缓存(默认):同一个A会和后续所有B配对,直到遇到下一个A替换缓存:
A1-B1 A1-B2 A2-B3 - 重置缓存(取消注释那行代码):每找到一对A+B就清空缓存,重新寻找下一个A:
A1-B1 A2-B3
关键细节说明
- Cold Flowable适配:因为
scan是在流的遍历过程中维护状态,每次订阅都会从头开始执行逻辑,完全适配Cold Flowable的特性; - 线程安全:RxJava的
scan操作符默认是串行执行的,只要没有额外指定多线程调度器,状态维护不会有线程安全问题; - 灵活性:你可以根据实际需求修改状态逻辑,比如支持多个A缓存、调整B的匹配规则等。
内容的提问来源于stack exchange,提问作者Namnodorel
相关产品推荐
相关产品推荐

