【问题标题】:How can Akka streams be materialized continually?Akka 流如何持续实现?
【发布时间】:2015-09-12 22:38:16
【问题描述】:

我在 Scala 中使用Akka Streams 使用AWS Java SDKAWS SQS 队列中进行轮询。我创建了一个ActorPublisher,它以两秒的间隔将消息出列:

class SQSSubscriber(name: String) extends ActorPublisher[Message] {
  implicit val materializer = ActorMaterializer()

  val schedule = context.system.scheduler.schedule(0 seconds, 2 seconds, self, "dequeue")

  val client = new AmazonSQSClient()
  client.setRegion(RegionUtils.getRegion("us-east-1"))
  val url = client.getQueueUrl(name).getQueueUrl

  val MaxBufferSize = 100
  var buf = Vector.empty[Message]

  override def receive: Receive = {
    case "dequeue" =>
      val messages = iterableAsScalaIterable(client.receiveMessage(new ReceiveMessageRequest(url).getMessages).toList
      messages.foreach(self ! _)
    case message: Message if buf.size == MaxBufferSize =>
      log.error("The buffer is full")
    case message: Message =>
      if (buf.isEmpty && totalDemand > 0)
        onNext(message)
      else {
        buf :+= message
        deliverBuf()
      }
    case Request(_) =>
      deliverBuf()
    case Cancel =>
      context.stop(self)
  }

  @tailrec final def deliverBuf(): Unit =
    if (totalDemand > 0) {
      if (totalDemand <= Int.MaxValue) {
        val (use, keep) = buf.splitAt(totalDemand.toInt)
        buf = keep
        use foreach onNext
      } else {
        val (use, keep) = buf.splitAt(Int.MaxValue)
        buf = keep
        use foreach onNext
        deliverBuf()
      }
    }
}

在我的应用程序中,我也尝试以 2 秒的间隔运行流程:

val system = ActorSystem("system")
val sqsSource = Source.actorPublisher[Message](SQSSubscriber.props("queue-name"))
val flow = Flow[Message]
  .map { elem => system.log.debug(s"${elem.getBody} (${elem.getMessageId})"); elem }
  .to(Sink.ignore)

system.scheduler.schedule(0 seconds, 2 seconds) {
  flow.runWith(sqsSource)(ActorMaterializer()(system))
}

但是,当我运行我的应用程序时,我会收到 java.util.concurrent.TimeoutException: Futures timed out after [20000 milliseconds] 以及随后由 ActorMaterializer 引起的死信通知。

是否有推荐的方法来持续实现 Akka 流?

【问题讨论】:

  • 我现在无法对其进行测试,但我不确定是否要使用多个 ActorMaterializer 实例。您在 ActorPublisher 中使用一个实例,而在整个流程中使用另一个实例。
  • 我最终使用了 Akka-Camel,因为它有一个很好的 SQS 集成,完成了我需要做的所有事情 (github.com/fzakaria/Akka-Camel-SQS)。
  • 您是否有理由必须每 2 秒连续使用不同的 ActorPublisher?根据给出的示例代码,继续使用同一个发布者会容易得多......

标签: scala akka amazon-sqs aws-sdk akka-stream


【解决方案1】:

我认为您不需要每 2 秒创建一个新的ActorPublisher。这似乎是多余的和浪费的内存。另外,我不认为 ActorPublisher 是必要的。根据我对代码的了解,您的实现将有越来越多的 Streams 都在查询相同的数据。来自客户端的每个Message 都会被 N 个不同的 akka Streams 处理,更糟糕的是,N 会随着时间的推移而增长。

无限循环查询的迭代器

您可以使用 scala 的 Iterator 从您的 ActorPublisher 获得相同的行为。可以创建一个不断查询客户端的迭代器:

//setup the client
val client = {
  val sqsClient = new AmazonSQSClient()
  sqsClient setRegion (RegionUtils getRegion "us-east-1")
  sqsClient
}

val url = client.getQueueUrl(name).getQueueUrl

//single query
def queryClientForMessages : Iterable[Message] = iterableAsScalaIterable {
  client receiveMessage (new ReceiveMessageRequest(url).getMessages)
}

def messageListIteartor : Iterator[Iterable[Message]] = 
  Iterator continually messageListStream

//messages one-at-a-time "on demand", no timer pushing you around
def messageIterator() : Iterator[Message] = messageListIterator flatMap identity

这个实现只在之前的所有消息都被消费完后才查询客户端,因此是真正的reactive。无需跟踪固定大小的缓冲区。您的解决方案需要一个缓冲区,因为消息的创建(通过计时器)与消息的消耗(通过 println)分离。在我的实现中,创建和消费是通过背压tightly coupled

Akka 流源

然后您可以使用此迭代器生成器函数来提供 akka 流源:

def messageSource : Source[Message, _] = Source fromIterator messageIterator

流动形成

最后,您可以使用此 Source 来执行 println(附带说明:您的 flow 值实际上是 Sink,因为 Flow + Sink = Sink)。使用问题中的 flow 值:

messageSource runWith flow

一个 akka Stream 处理所有消息。

【讨论】:

    猜你喜欢
    • 2016-04-09
    • 1970-01-01
    • 1970-01-01
    • 2017-11-16
    • 2016-09-26
    • 1970-01-01
    • 2013-04-02
    • 2012-07-11
    • 2017-01-30
    相关资源
    最近更新 更多