如何让一个Observable等待另一个,且后者获取其最后发射值?
RxJava 两个需求的解决方案
咱们逐个拆解你的技术需求,都是RxJava里典型的流组合场景,直接上实用的实现方案:
需求1:让一个Observable等待另一个Observable,且后者仅获取前者最后发射的值
这里我先明确两种常见的场景(你描述的需求大概率对应其中一种),分别给出实现:
场景1:等待的Observable(B)仅在源Observable(A)至少发射过一次值后才启动,且每次B发射都用A的最新值
这种情况用withLatestFrom操作符最适合,但要注意:如果A还没发射值,B的初始发射会被直接忽略。如果要强制B必须等A的第一个值出来再启动,可以结合publish避免重复订阅A:
// 替换成你自己的数据流即可 Observable<Integer> observableA = Observable.just(1, 2, 3).delay(1, TimeUnit.SECONDS); Observable<String> observableB = Observable.just("x", "y").delay(2, TimeUnit.SECONDS); observableA.publish(aStream -> // 先等待A发射第一个值,再启动B的处理逻辑 aStream.take(1) .ignoreElements() .andThen(observableB.withLatestFrom(aStream, (bVal, aVal) -> String.format("B的内容: %s,A的最新值: %d", bVal, aVal) )) ).subscribe(System.out::println);
运行后你会看到,B的发射会严格等到A至少有一个值后才开始,并且每次B发射都会取A当前的最后一个值。
场景2:等待的Observable(B)要等源Observable(A)完全完成,再获取A的最后一个值继续处理
如果A是有限流(会最终完成),且你只需要A的最后一个值,用takeLast(1)结合concatMap就能实现:
Observable<Integer> observableA = Observable.just(1, 2, 3).delay(1, TimeUnit.SECONDS); Observable<String> observableB = Observable.just("x", "y"); observableA.takeLast(1) .concatMap(lastAVal -> observableB.map(bVal -> String.format("B的内容: %s,A的最后值: %d", bVal, lastAVal) ) ) .subscribe(System.out::println);
这里A完全完成后,B才会开始处理,所有B的发射都会复用A的最后一个值。
需求2:结合两个Observable更新列表,首次需等待两者都发射第一批数据
这个场景的核心是按ID匹配列表项和状态,还要确保首次处理必须等A和B都有初始数据。我们可以用scan维护两个缓存Map,再用combineLatest合并生成完整的带状态列表:
首先定义你的数据结构(替换成你实际的业务类即可):
// 列表项实体 class ListItem { String id; String content; public ListItem(String id, String content) { this.id = id; this.content = content; } } // 状态实体(仅含ID和状态标记) class ItemStatus { String id; boolean isActive; public ItemStatus(String id, boolean isActive) { this.id = id; this.isActive = isActive; } }
然后实现流的组合逻辑:
// 模拟你的Observable A:发射列表项 Observable<ListItem> observableA = Observable.just( new ListItem("a", "项目A"), new ListItem("b", "项目B"), new ListItem("c", "项目C") ).delay(1, TimeUnit.SECONDS); // 模拟你的Observable B:发射对应ID的状态 Observable<ItemStatus> observableB = Observable.just( new ItemStatus("a", true), new ItemStatus("b", false), new ItemStatus("c", true) ).delay(1500, TimeUnit.MILLISECONDS); // 模拟B的初始发射稍晚于A // 用scan维护列表项的缓存Map:每次A发射新项,更新缓存 Observable<Map<String, ListItem>> itemCache$ = observableA.scan(new HashMap<>(), (cache, item) -> { cache.put(item.id, item); return new HashMap<>(cache); // 返回新Map,确保下游能感知到更新 }).distinctUntilChanged(); // 避免重复发射相同的缓存 // 用scan维护状态的缓存Map:每次B发射新状态,更新缓存 Observable<Map<String, Boolean>> statusCache$ = observableB.scan(new HashMap<>(), (cache, status) -> { cache.put(status.id, status.isActive); return new HashMap<>(cache); }).distinctUntilChanged(); // 合并两个缓存,生成已匹配的带状态列表 Observable.combineLatest(itemCache$, statusCache$, (itemMap, statusMap) -> { // 筛选出同时存在于两个缓存中的项(ID匹配) return itemMap.entrySet().stream() .filter(entry -> statusMap.containsKey(entry.getKey())) .map(entry -> { ListItem item = entry.getValue(); boolean status = statusMap.get(entry.getKey()); return String.format("ID: %s | 内容: %s | 状态: %s", item.id, item.content, status ? "激活" : "未激活"); }) .collect(Collectors.toList()); }) // 过滤空列表,确保首次发射是A和B都有初始数据之后的结果 .filter(list -> !list.isEmpty()) .subscribe(updatedList -> { System.out.println("=== 列表已更新 ==="); updatedList.forEach(System.out::println); });
这段代码的核心逻辑:
scan操作符持续维护缓存Map,每次A或B有新数据就更新缓存,返回新Map对象触发下游更新。combineLatest会在任意一个缓存更新时,合并两者的数据,筛选出ID匹配的项和状态,生成完整列表。filter(list -> !list.isEmpty())确保首次发射必须等A和B都至少发射了一条数据(此时两个缓存都有内容,匹配后的列表不为空)。- 后续不管A发射新项还是B发射新状态,都会自动更新列表,保持数据一致。
内容的提问来源于stack exchange,提问作者Cilvet
相关产品推荐
相关产品推荐

