【问题标题】:Consumer/Producer AWS SQS akka scala with synchrone consumer消费者/生产者 AWS SQS akka scala 与同步消费者
【发布时间】:2016-06-17 13:06:31
【问题描述】:

我的应用程序有一个生产者和一个消费者。我的制作人不定期地生成消息。有时我的队列会是空的,有时我会有一些消息。 我想让我的消费者收听队列,当有消息时,接受它并处理这个消息。这个过程可能需要几个小时,如果我的消费者没有完成对当前消息的处理,我不希望我的消费者从队列中获取另一条消息。

我认为 AKKA 和 AWS SQS 可以满足我的需求。通过阅读文档和示例,akka-camel 似乎可以简化我的工作。

我在 github 上找到了这个 example

我对消费者的配置更感兴趣(Thread.sleep只是为了模拟我的处理):

class MySqsConsumer extends Consumer {

  //The SQS URI is an in-only message exchange (autoAck=true)
  override def endpointUri = "aws-sqs://sqs-akka-camel?  amazonSQSClient=#client"

  override def receive = {
    case msg: CamelMessage => {
      println("Start received %s" format msg.bodyAs[String])
      Thread.sleep(4000)
      println("Stop received %s" format msg.bodyAs[String])

    }

  }
 }

堆栈跟踪:

...
14:38:53.335 [Camel (sqs-akka-camel) thread #0 - aws-sqs://sqs-akka-camel] DEBUG o.a.camel.processor.SendProcessor - >>>> Endpoint[akka://sqs-akka-camel/user/consumer?autoAck=true&replyTimeout=60000+milliseconds] Exchange[Message: Hello World1!]
14:38:54.051 [Camel (sqs-akka-camel) thread #0 - aws-sqs://sqs-akka-camel] DEBUG o.a.camel.processor.SendProcessor - >>>> Endpoint[akka://sqs-akka-camel/user/consumer?autoAck=true&replyTimeout=60000+milliseconds] Exchange[Message: Hello World2!]
14:38:54.753 [Camel (sqs-akka-camel) thread #0 - aws-sqs://sqs-akka-camel] DEBUG o.a.camel.processor.SendProcessor - >>>> Endpoint[akka://sqs-akka-camel/user/consumer?autoAck=true&replyTimeout=60000+milliseconds] Exchange[Message: Hello World3!]
Start received Hello World1!
Stop received Hello World1!
Start received Hello World2!
Stop received Hello World2!
Start received Hello World3!
...

我的问题是 akka-camel 在接收方法完成之前从队列中读取消息。

如何让我的消费者在接收新消息之前等待流程结束?如果我使用了错误的工具/库,您能否指导我使用新工具?

【问题讨论】:

    标签: scala akka amazon-sqs producer-consumer akka-camel


    【解决方案1】:

    不确定您是否使用 akka 流,但如果使用,请将源发送到 Sink.queue

    val graph=yourSource.to(Sink.queue)
    

    当物化(graph.run)时,它会返回一个 SinkQueue,你可以手动从中拉取项目。它将使用背压机制,因此您不必担心不拉取物品会发生什么。

    注意:我故意忽略了泛型!但不要忘记它们。

    【讨论】:

      【解决方案2】:

      据我所知,你做不到。 Camel 一直在消费队列中的消息。

      您可以设置override def autoAck = false - 在您确认当前消息之前,它不会向receive() 发送另一条消息。但是“隐藏”的骆驼演员仍然会消耗。

      也许你可以提供一个自定义的amazonSQSClient 并以某种方式控制它。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2018-10-16
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2019-08-19
        • 2018-12-13
        • 2018-01-07
        相关资源
        最近更新 更多