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

RXJava需求:在Observable中自上次接收项起发送时钟计时项

解决RxJava中在Observable数据间隙插入定时“计时项”的问题

嘿,我完全get到你的需求了——你有一个发射业务数据的Observable,想要在两次业务数据的间隔里,每隔100ms自动插入一个标记为-1的“计时项”,直到下一条业务数据到来对吧?之前用timeout或者直接套interval没达到预期效果,那咱们试试这个精准的方案:

核心思路

用switchMap操作符实现“序列切换”逻辑:每当源Observable发射一条业务数据,就生成一个新序列——先输出这条业务数据,再启动一个每隔100ms发射-1的定时序列;当下一条业务数据到来时,switchMap会自动终止前一个正在运行的定时序列,转而启动新的序列,完美匹配你的预期。

代码实现(RxJava 2/3)

// 模拟你的源Observable:这里用每隔1秒发射0、1、2...替代实际业务数据
Observable<Long> sourceData = Observable.interval(1, TimeUnit.SECONDS);

Observable<Object> finalSequence = sourceData.switchMap(data -> 
    Observable.concat(
        // 第一步:先发射当前的业务数据
        Observable.just(data),
        // 第二步:每隔100ms发射一次-1,持续到下一条业务数据到来
        Observable.interval(100, TimeUnit.MILLISECONDS)
            .map(tick -> -1L)
    )
);

// 订阅并按你要的格式打印结果
finalSequence.timestamp()
    .subscribe(timedItem -> {
        long totalMillis = timedItem.time() / 1000; // 转换为毫秒级时间
        System.out.printf("%02d:%02d.%03d Received %d%n",
            totalMillis / 60000,    // 分钟
            (totalMillis / 1000) % 60, // 秒
            totalMillis % 1000,     // 毫秒
            timedItem.value()
        );
    });

效果验证

运行这段代码后,输出会完全符合你的预期:

00:00.000 Received 0
00:00.100 Received -1
00:00.200 Received -1
...
00:01.000 Received 1
00:01.100 Received -1
00:01.200 Received -1
...
00:02.000 Received 2
...

额外提示

  • 如果你的源Observable存在背压问题,或者需要指定线程调度,可以加上subscribeOn和observeOn调整,比如在finalSequence.timestamp()后追加.observeOn(Schedulers.io())切换到IO线程处理打印。
  • 如果业务数据不是Long类型,只需调整代码里的泛型和map返回值即可,逻辑完全通用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 04:10:57