【问题标题】:Alpakka AMQP : How to detect declaration exception?Alpakka AMQP:如何检测声明异常?
【发布时间】:2020-05-12 20:55:44
【问题描述】:

我有一个带有声明的 AMQP 源和 AMQP 接收器:

List<Declaration> declarations = new ArrayList<Declaration>() {{
      add(QueueDeclaration.create(sourceExchangeName));
      add(BindingDeclaration.create(sourceExchangeName, sourceExchangeName).withRoutingKey(sourceRoutingKey));
    }};

amqpSource = AmqpSource
    .committableSource(
        NamedQueueSourceSettings.create(connectionProvider, sourceExchangeName)
            .withDeclarations(declarations),
        bufferSize);

AmqpWriteSettings amqpWriteSettings = AmqpWriteSettings.create(connectionProvider)
    .withExchange("DEST_XCHANGE")
    .withRoutingKey("ROUTE123")
    .withDeclaration(ExchangeDeclaration.create(destinationExchangeName,
        BuiltinExchangeType.DIRECT.getType()));

amqpSink = AmqpSink.create(amqpWriteSettings);

然后我有一个流程..

amqpSource.map(doSomething).async().map(doSomethingElse).async().to(amqpSink)

现在,在我启动应用程序后,发送到源队列的消息没有被消耗。后来我发现这是由于声明期间发生的错误。 (即,当我在 Source 和 Sink 设置中删除 .withDeclarations(..) 时,它运行良好。

所以我的问题:

  1. 如何检测 AMQP Source 和 Sink 是否正常运行?
  2. 如何忽略声明异常?
  3. 如果出现异常,如何知道并导致系统出现故障?

【问题讨论】:

    标签: akka akka-stream akka-http akka-cluster alpakka


    【解决方案1】:

    要回答 1 和 3,AmqpSink 实现了一个 CompletionStage&lt;Done&gt;,您必须保留并处理(注册一些回调函数)以观察流的失败和完成。在文档示例中,我们阻止了在生产代码中不好的完成阶段 (https://doc.akka.io/docs/alpakka/current/amqp.html#with-sink),这可能是因为该示例包含在 Alpakka 测试之一中。更喜欢通常的CompletionStage 回调/转换方法(例如参见this introduction)。

    CompletionStage 将在发生错误时失败,当流被物化/启动时或在元素处理期间,或者在源到达末尾并且每个元素都通过您的流进入接收器时完成。这意味着要启动流,如果它没有很快失败,它就会运行。

    对于问题 2,不确定是否可以忽略声明异常,可能是那些总是连接失败。

    【讨论】:

    • 当我使用 AmqpFlow 进行写作,然后进行更多映射时......最后我做了一个 Sink.Ignore。在这种情况下如何检测故障?
    • 流程中的故障会顺流而下并最终出现在Sink.ignore 的具体化值中,因此请务必坚持下去。
    • 谢谢。所以这意味着当蒸汽的完成阶段完成时,流已经完成并且不再运行。谢谢!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2014-08-02
    • 2014-05-06
    • 2013-03-22
    • 2023-04-01
    • 1970-01-01
    • 1970-01-01
    • 2011-09-26
    相关资源
    最近更新 更多