【发布时间】:2015-09-12 22:38:16
【问题描述】:
我在 Scala 中使用Akka Streams 使用AWS Java SDK 从AWS 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