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

如何修改Akka流Throttler实现Scala中消息间自定义间隔延迟?

实现Akka Stream中每条消息的可变发送间隔

当然可以实现自定义每条消息之间的发送间隔!你当前使用的throttle算子是针对固定速率限流的,没法直接实现不同消息间的可变延迟,但我们可以通过Akka Stream的组合算子来灵活实现这个需求。

核心思路

我们可以给每条消息绑定对应的延迟时间,然后通过流的顺序处理逻辑,让每条消息在上一条消息发送完成后,等待指定的延迟时间再发送到目标Actor。具体来说:

  1. 定义一个包含消息内容和延迟时长的包装类
  2. 使用flatMapConcat算子将每条带延迟的消息转换成一个带延迟的单元素流,保证顺序执行
  3. 最终将处理后的消息发送到目标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 ! _)
  }
}

关键改动说明

  1. 新增DelayedAction类:用来将消息内容和对应的延迟时间绑定,让我们能为每条消息指定不同的间隔
  2. 重构getThrottler方法:使用flatMapConcat替代原有的throttle算子,它会按顺序处理每个DelayedAction,对每条消息先执行延迟操作,再发送到目标Actor,确保消息顺序和延迟要求
  3. 调整消息发送逻辑:在main方法中构造带延迟的消息列表,明确指定每条消息与前一条的间隔时间
  4. 添加时间戳打印:方便你直观看到每条消息的发送时间,验证延迟是否符合预期

额外说明

  • 如果你的延迟规则需要根据消息内容动态计算(比如根据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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 05:27:28