如何在Akka Streams中限制Source触发频率并仅处理最新消息?
问题分析
你遇到的核心问题是:Akka Stream的throttle是先获取元素再等待时间窗口,导致它拿到buffer里的旧元素后,即使buffer后续被新元素替换,也只会输出之前拿到的旧值,无法满足“到时间点处理最新事件”的需求。
你的需求本质是:
- 固定间隔(7秒)处理一次事件
- 间隔内的新事件完全替换旧事件,只保留最新的那个
- 间隔结束时,处理当前最新的事件
解决方案
不需要自定义Flow,用Akka Stream原生操作符就能解决,以下两种方案都能满足需求:
方案一:conflate + throttle(推荐)
Source.tick(5.seconds, 3.seconds, 1) .scan(0)((a, b) => a + b) // 生成递增计数器 .wireTap(num => logger.warn(s"up ${num.formatted("%02d")}")) // 合并上游所有新元素,只保留最新的那个 .conflate((_, newElem) => newElem) // 每隔7秒输出一个元素,Shaping模式确保严格按速率输出 .throttle(1, 7.seconds, 1, ThrottleMode.Shaping) .wireTap(num => logger.warn(s"down ${num.formatted("%02d")}")) .runWith(Sink.ignore)(materializer)
原理说明:
conflate:当下游处理速度慢于上游时,会把上游的多个元素合并成一个。这里我们定义合并逻辑为保留最新元素,所以上游不管发多少事件,内存里始终只存当前最新的那个。throttle(Shaping模式):每隔7秒向conflate请求一次元素,此时conflate会返回当前最新的事件,完美匹配你的需求。
方案二:定时器 + zipLatest
如果需要更直观的触发逻辑,可以用定时器作为触发源,和事件流做zipLatest:
// 原事件源,用conflate保留最新元素 val eventSource = Source.tick(5.seconds, 3.seconds, 1) .scan(0)((a, b) => a + b) .wireTap(num => logger.warn(s"up ${num.formatted("%02d")}")) .conflate((_, newElem) => newElem) // 每隔7秒触发一次,拉取最新事件 Source.tick(0.seconds, 7.seconds, ()) .zipLatest(eventSource) .map(_._2) // 只取事件源的最新元素 .wireTap(num => logger.warn(s"down ${num.formatted("%02d")}")) .runWith(Sink.ignore)(materializer)
原理说明:
Source.tick生成固定间隔的触发信号zipLatest会在触发信号到达时,取事件源的当前最新元素,这样每次触发都是处理最新的事件
验证效果
用上述方案运行后,日志会变成类似这样:
up 01 down 01 up 02 up 03 down 03 up 04 up 05 up 06 down 06 up 07 up 08 down 08 up 09 up 10 down 10
可以看到每次down输出的都是间隔内最后一个up的最新元素,符合预期。
内容的提问来源于stack exchange,提问作者Uko
相关产品推荐
相关产品推荐

