【问题标题】:Infinite AMQP Consumer with Alpakka使用 Alpakka 的无限 AMQP 消费者
【发布时间】:2019-03-15 00:23:25
【问题描述】:

我正在尝试使用 Alpakka 实现一个连接到 AMQP 代理的非常简单的服务。我只是希望它在将消息推送到给定的交换/主题时将其队列中的消息作为流使用。

在我的测试中似乎一切正常,但是当我尝试启动我的服务时,我意识到我的流只使用了我的消息一次然后退出。

基本上我使用的是 Alpakka 文档中的代码:

def consume()={
    val amqpSource = AmqpSource.committableSource(
      TemporaryQueueSourceSettings(connectionProvider, exchangeName)
        .withDeclaration(exchangeDeclaration)
        .withRoutingKey(topic),
      bufferSize = prefetchCount
    )

    val amqpSink = AmqpSink.replyTo(AmqpReplyToSinkSettings(connectionProvider))

    amqpSource.mapAsync(4)(msg => onMessage(msg)).runWith(amqpSink)
}

我尝试安排每秒执行一次 consume(),但遇到了 OutOfMemoryException 问题。

有什么合适的方法可以让这段代码无限循环运行吗?

【问题讨论】:

    标签: akka-stream alpakka


    【解决方案1】:

    如果您想让Source 在失败或被取消时重新启动,请使用RestartSource.withBackoff 包装它。

    【讨论】:

    • 我已经尝试过使用RestartSource.withBackoff,但是当它完成时并没有重新启动源,只有当它失败时。我正在寻找一个永无止境的来源。
    • onFailuresWithBackoff 这样做。 withBackoff 在失败和完成时都重新启动。
    • 你说得对,RestartSource 做到了,谢谢!我对RestartSource 的问题来自我的接收器。现在似乎一切正常。
    猜你喜欢
    • 1970-01-01
    • 2020-07-15
    • 2011-09-08
    • 2023-03-16
    • 1970-01-01
    • 2019-09-19
    • 1970-01-01
    • 1970-01-01
    • 2020-02-22
    相关资源
    最近更新 更多