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

如何将Observable发射的列表与另一个Observable执行zip操作?

RxJava:定时Observable与序列Observable的动态Zip实现

需求说明

  • Observable 1:每1秒发射一个递增数字(1、2、3……)
  • Observable 2:每3秒发射一组数字序列(例如[10,20,30,40])
  • 期望效果:每次Observable 2发射新序列时,将序列中的元素依次与Observable 1后续发射的元素配对;旧序列未匹配完时,若新序列到来则丢弃旧序列剩余元素,改用新序列继续配对;无序列时直接输出Observable 1的元素。

效果示例

Observable 1Observable 2zip结果
1[1]
2[2]
3[10,20,30,40][3, 10]
4[4, 20]
5[5, 30]
6[10,20,30,40][6, 40, 10]
7[7, 20]

实现代码(RxJava)

import io.reactivex.rxjava3.core.Observable;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.TimeUnit;

public class DynamicZipExample {
    public static void main(String[] args) throws InterruptedException {
        // Observable 1:每秒发射递增数字(1、2、3...)
        Observable<Long> obs1 = Observable.interval(1, TimeUnit.SECONDS)
                .map(num -> num + 1);

        // Observable 2:每3秒发射一次数字序列
        Observable<List<Integer>> obs2 = Observable.interval(3, TimeUnit.SECONDS)
                .map(num -> Arrays.asList(10, 20, 30, 40));

        // 核心逻辑:动态切换待匹配序列,完成按需zip
        obs1.switchMap(currentObs1Num -> obs2
                .startWith(Collections.emptyList()) // 初始阶段提供空序列,处理无匹配项的情况
                .switchMap(sequence -> {
                    if (sequence.isEmpty()) {
                        // 无待匹配序列时,直接返回当前Observable 1的元素包装列表
                        return Observable.just(Collections.singletonList(currentObs1Num));
                    } else {
                        // 将序列转为Observable,与当前及后续Observable 1元素配对
                        return Observable.fromIterable(sequence)
                                .zipWith(Observable.just(currentObs1Num)
                                        .concatWith(obs1.skipUntil(Observable.timer(1, TimeUnit.SECONDS))),
                                        (seqItem, num) -> Arrays.asList(num, seqItem))
                                .take(1); // 每次仅取当前Observable 1元素对应的配对结果
                    }
                })
                .take(1))
                .subscribe(result -> System.out.println("输出结果:" + result));

        // 保持程序运行,观察输出
        Thread.sleep(10000);
    }
}

代码逻辑解析

  1. Observable初始化:分别创建两个定时发射的Observable,obs1每秒生成递增数字,obs2每3秒生成固定序列。
  2. 动态序列切换:使用switchMap监听obs1的每个元素,同时通过obs2.startWith(emptyList)处理初始无序列的场景。
  3. 配对逻辑:
    • 无序列时直接输出obs1的元素;
    • 有序列时将序列转为Observable,与obs1的当前及后续元素进行zip配对;
    • 当obs2发射新序列时,switchMap会自动丢弃旧序列的未匹配元素,切换为新序列继续配对,符合示例中的效果。

内容的提问来源于stack exchange,提问作者allen kim

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 19:10:19