如何修改Akka流Throttler实现Scala中消息间自定义间隔延迟?
实现Akka Stream中每条消息的可变发送间隔
当然可以实现自定义每条消息之间的发送间隔!你当前使用的throttle算子是针对固定速率限流的,没法直接实现不同消息间的可变延迟,但我们可以通过Akka Stream的组合算子来灵活实现这个需求。
核心思路
我们可以给每条消息绑定对应的延迟时间,然后通过流的顺序处理逻辑,让每条消息在上一条消息发送完成后,等待指定的延迟时间再发送到目标Actor。具体来说:
- 定义一个包含消息内容和延迟时长的包装类
- 使用
flatMapConcat算子将每条带延迟的消息转换成一个带延迟的单元素流,保证顺序执行 - 最终将处理后的消息发送到目标Actor
修改后的完整代码
import akka.NotUsed import akka.actor.{Actor, ActorRef, ActorSystem, Props} import akka.stream.{ActorMaterializer, OverflowStrategy} import akka.stream.scaladsl.{Sink, Source} import scala.concurrent.duration.{FiniteDuration, _} // 原有消息类 case class Action(side: String) // 新增:绑定消息与延迟时间的包装类 case class DelayedAction(action: String, delay: FiniteDuration) object PlaceTrade { val actorSystem = ActorSystem("firstActorSystem") println(actorSystem.name) object TradeAction { def props(action : String) = Props(new TradeAction(action)) } class TradeAction(actorName: String) extends Actor { override def receive: Receive = { case "Buy" => { val r = requests.get("http://www.google.com") println(s"[${java.time.LocalTime.now()}] r status code is ${r.statusCode}") println(s"[${java.time.LocalTime.now()}] Buy") println("") } case "Sell" => { val r = requests.get("http://www.google.com") println(s"[${java.time.LocalTime.now()}] r status code is ${r.statusCode}") println(s"[${java.time.LocalTime.now()}] Sell") println("") } case _ => } } implicit val materializer = ActorMaterializer.create(actorSystem) // 修改后的限流/延迟处理逻辑 def getThrottler(ac: ActorRef) = Source.actorRef[DelayedAction](bufferSize = 1000, OverflowStrategy.dropNew) .flatMapConcat { delayed => // 对每条消息,等待指定延迟后再发送 Source.single(delayed.action).delay(delayed.delay) } .to(Sink.actorRef(ac, NotUsed)) .run() def main(args: Array[String]): Unit = { val tradeAction = actorSystem.actorOf(TradeAction.props("TradeAction")) val throttler = getThrottler(tradeAction) // 构造带自定义延迟的消息列表: // 第一条立即发送,第二条与第一条间隔2秒,第三条与第二条间隔1秒,第四条与第三条间隔3秒 val delayedActions = List( DelayedAction("Buy", 0.second), DelayedAction("Buy", 2.second), DelayedAction("Buy", 1.second), DelayedAction("Sell", 3.second) ) delayedActions.foreach(throttler ! _) } }
关键改动说明
- 新增
DelayedAction类:用来将消息内容和对应的延迟时间绑定,让我们能为每条消息指定不同的间隔 - 重构
getThrottler方法:使用flatMapConcat替代原有的throttle算子,它会按顺序处理每个DelayedAction,对每条消息先执行延迟操作,再发送到目标Actor,确保消息顺序和延迟要求 - 调整消息发送逻辑:在
main方法中构造带延迟的消息列表,明确指定每条消息与前一条的间隔时间 - 添加时间戳打印:方便你直观看到每条消息的发送时间,验证延迟是否符合预期
额外说明
- 如果你的延迟规则需要根据消息内容动态计算(比如根据Action的side来决定延迟),可以在
flatMapConcat中加入计算逻辑,比如:.flatMapConcat { delayed => val calculatedDelay = delayed.action match { case "Buy" => 2.second case "Sell" => 1.second case _ => 0.second } Source.single(delayed.action).delay(calculatedDelay) } flatMapConcat保证了消息的顺序性,不会出现乱序发送的情况,这对于交易类场景非常重要
内容的提问来源于stack exchange,提问作者blue-sky
相关产品推荐
相关产品推荐

