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

RxJava需求:周期性发送上游已接收的最新数据项

RxJava 周期性发送上游最新数据项的解决方案

你需要的是能周期性触发并发送上游流截至当前时刻的最新数据的逻辑,throttleLatest仅保留限流间隔内的最新项,确实不符合需求。下面提供两种实现方式:

方式一:现有操作符组合实现

利用publish共享上游流,结合withLatestFrom让周期流每次触发时获取上游最新值:

// 定义上游流,用publish实现多订阅共享
val upstream = Flowable.interval(3, TimeUnit.SECONDS)
    .publish()

// 每1秒触发一次,获取上游最新数据并发送
Flowable.interval(1, TimeUnit.SECONDS)
    .withLatestFrom(upstream) { _, latestItem -> latestItem }
    .subscribe { println(it) }

// 连接共享流,启动上游发射
upstream.connect()

执行逻辑说明:

  • 上游每3秒发射递增数字(0、1、2...)
  • 下游周期流每1秒触发一次,每次取上游当前的最新值发送
  • 最终输出为0, 0, 0, 1, 1, 1, 2, 2, 2...,完全匹配你的期望。

方式二:封装自定义操作符

如果需要复用该逻辑,可以封装成自定义操作符,和你示例中的operatorThatEmitsMostRecentItemFromUpstream对应:

fun <T> Flowable<T>.emitMostRecentPeriodically(period: Long, unit: TimeUnit): Flowable<T> {
    return this.publish { upstream ->
        Flowable.interval(period, unit)
            .withLatestFrom(upstream) { _, latest -> latest }
    }
}

// 调用自定义操作符
Flowable.interval(3, TimeUnit.SECONDS)
    .emitMostRecentPeriodically(1, TimeUnit.SECONDS)
    .subscribe { println(it) }

该操作符内部复用了方式一的逻辑,对外提供更简洁的调用方式,和你示例代码的写法完全匹配。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 09:54:16