【问题标题】:How to programmatically terminate io.scalac.amqp.Connection from reactive-rabbit library如何以编程方式从反应兔库中终止 io.scalac.amqp.Connection
【发布时间】:2015-05-14 04:34:21
【问题描述】:

我正在结合使用 akka 流和响应式兔子库来构建一个脚本,该脚本将一些信息推送到我本地 rabbitmq 服务器上的交换。

一旦信息被推送到队列中,我希望程序自行关闭。但是Connection 使程序保持活力,我在Connection 上找不到任何方法或其他如何杀死它的示例。不可避免地,我必须手动终止该进程。

我的代码如下所示:

package prototype

import akka.actor.ActorSystem
import akka.stream.FlowMaterializer
import akka.stream.scaladsl.{Sink, Source}
import akka.util.ByteString
import io.scalac.amqp.{Message, Connection}

object PopulateTodoQueue extends App {
  val connection = Connection()

  val message = Message(ByteString("message"))

  val source = Source(List(message))
  val sink = Sink(connection.publishDirectly(queue = "todo"))

  implicit val actorSystem = ActorSystem()
  implicit val materializer = FlowMaterializer()

  (source to sink).run()

  // Quick hack to wait long enough for the message to send
  Thread.sleep(1000)
  actorSystem.shutdown()
}

这是来自我的 build.sbt 库依赖项的 sn-p:

"com.typesafe.akka"          %%  "akka-actor"               % "2.3.7",
"com.typesafe.akka"          %%  "akka-stream-experimental" % "0.11",
"io.scalac"                  %%  "reactive-rabbit"          % "0.2.1",

对于这些一次性任务是否有更好的模式 - 例如您将回调传递给的临时连接?我见过的示例中的所有用例都是针对长时间运行的客户端,这些客户端一直运行到用户明确杀死它们为止。

提前致谢!

【问题讨论】:

  • +1,我目前在想同样的事情,虽然我的用例是相反的,使用反应兔订阅...
  • @Tycho 这个问题已经在这里很久了,所以我认为你不会得到答案。我不知道他们甚至不再维护图书馆。我最终只是回到了简单的 java 驱动程序并为它编写了一个小包装器来抽象消费者的阻塞性质。您可以给公司发电子邮件 (scalac.io)。我在 scala 驱动程序中遇到过很多此类问题 - 它们通常过于复杂或维护/支持不佳。

标签: scala rabbitmq reactive-programming akka-stream


【解决方案1】:

连接特征中有关闭方法

/** Shutdowns underlying connection.
    * Publishers and subscribers are terminated and notified via `onError`.
    * This method waits for all close operations to complete. */
  def shutdown(): Future[Unit]

你可以这样称呼

connection.shutdown()

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-11-20
    • 1970-01-01
    • 2022-10-21
    • 2022-11-12
    • 2015-05-31
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多