【问题标题】:Apache camel - how to "wiretap" synchronously? Or just send a copy of an exchange?Apache camel - 如何同步“窃听”?或者只是发送一份交换副本?
【发布时间】:2014-01-07 10:34:47
【问题描述】:

我有一个 apache 骆驼路由,它正在交换体上处理 POJO。

请查看从 1 到 3 标记的行的顺序。

    from("direct:foo")
        .to("direct:doSomething")         // 1 (POJO on the exchange body)
        .to("direct:storeInHazelcast")    // 2 (destroys my pojo! it gets -1)
        .to("direct:doSomethingElse")     // 3 (Where is my POJO??)
    ;

现在我需要在hazelcast 组件上使用put 操作,不幸的是它需要将主体设置为值-1。

    from("direct:storeInHazelcast")
            .setBody(constant(-1))
            .setHeader(HazelcastConstants.OPERATION, constant(HazelcastConstants.PUT_OPERATION))
            .setHeader(HazelcastConstants.OBJECT_ID, constant(LAST_FLIGHT_UPDATE_SEQ))
            .to("hazelcast:map:MyNumber")
    ;

对于标记为 2 的行,我想将交换的副本发送到 storeInHazelcast 路由。

首先,我尝试了.multicast(),但交换体仍然搞砸了(到-1)。

        // shouldnt this copy the exchange?
        .multicast().to("direct:storeInHazelcast").end()

然后我尝试了.wireTap(),它作为“即发即弃”(异步)模式工作,但我实际上需要它来阻止,并等待它完成。可以做wireTap拦截吗?

        // this works but I need it to be sync processing (not async)
        .wireTap("direct:storeInHazelcast").end()

所以我在这里寻找一些提示。据我所知,multicast() 应该复制了交换,但我的 storeInHazelcast 路由中的 setBody() 看起来搞砸了原始交换。

或者,也许还有其他方法可以做到这一点。

提前致谢。 骆驼2.10

【问题讨论】:

  • 多播做浅拷贝,这反过来也反映在其他多播路由中。

标签: java apache-camel hazelcast jbossfuse


【解决方案1】:

我想我偶然发现了答案,2 行可以像这样使用来自 dsl 的enrich()

    .enrich("direct:storeInHazelcast", new KeepOriginalAggregationStrategy())

地点:

public class KeepOriginalAggregationStrategy implements AggregationStrategy {
    @Override
    public Exchange aggregate(Exchange oldExchange, Exchange newExchange) {
        return oldExchange;
    }
}

有趣的是,我找到了一个名为UseOriginalAggregationStrategy()的聚合策略,但是我看不到如何从DSL中指定名为Exchange original的参数。

    .enrich("direct:storeInHazelcast",
        new UseOriginalAggregationStrategy(???, false))

由于 dsl 中没有某种 getExchange() 方法,我在这里看不到如何使用此聚合策略(但如果有人可以建议如何使用,请执行)。

【讨论】:

  • 虽然这个答案有效,但我不同意使用enrich,因为它带有不同的语义并暗示了不同的模式。
  • @AlessandroDaRugna 我同意,但我还没有看到更好的解决方案。罗伯特的解决方案还可以,但我认为没有必要。 Camel 是关于集成模式的,真的没有适合这个用例的实现模式吗?
  • enrich 仅意味着产生一个 Exchange - 返回的 Exchange 发生的事情被封装在 AggregationStrategy 中。请参阅 Arek Bazylewicz 的答案以获得一个不错的解决方案(我也不知道有预定义的策略):.enrich("direct:storeInHazelcast", AggregationStrategies.useOriginal())
【解决方案2】:

您无需编写自己的聚合策略即可使用

.enrich("direct:storeInHazelcast", AggregationStrategies.useOriginal())

【讨论】:

    【解决方案3】:

    将其保存在标题中并恢复它。

    from("direct:foo")
        .to("direct:doSomething")         // 1 (POJO on the exchange body)
        .setHeader("old_body", body())    // save body
        .to("direct:storeInHazelcast")    // 2 (destroys my pojo! it gets -1)
        .setBody(header("old_body"))      // RESTORE the body
        .removeHeader("old_body")         // cleanup header
        .to("direct:doSomethingElse")     // 3 (Where is my POJO??)
    ;
    

    这是破坏性组件相当常见的范例。

    【讨论】:

    • 谢谢,我希望避免整个“临时变量”的事情,但这将责任转移到我的路线上,而不是要求调用者使用丰富。
    • 实际上,我想我应该在storeInHazelcast 路由内对标头执行临时变量,这将使我的 1 - 2 - 3 行代码保持不变。
    【解决方案4】:

    我也有这个要求(在另一条路由上执行同步、仅处理),为了实现它,我编写了一个自定义处理器,它以编程方式发送 Exchange 的副本。我认为这会产生更好的 DSL,其中使用点的语义比使用丰富更清晰。

    这个静态辅助方法创建处理器:

    public static Processor synchronousWireTap(String uri) {
        return exchange -> {
            Exchange copy = exchange.copy();
    
            exchange.getContext().createProducerTemplate().send(uri,copy);
    
            //ProducerTemplate.send(String,Exchange) does not, unlike other send methods, rethrow an exception
            //on the exchange. We want any unhandled exception to be rethrown, so we must do so here.
            Throwable thrown = copy.getException(Throwable.class);
    
            if (thrown != null) {
                throw new CamelExecutionException(thrown.getMessage(), exchange, thrown);
            }
        };
    }
    

    下面是一个使用示例:

    from("direct:foo")
        .to("direct:doSomething")                               // 1 (POJO on the exchange body)
        .process(synchronousWireTap("direct:storeInHazelcast")) // 2 (Does not destroy POJO because a copy of the exchange gets sent to this uri)
        .to("direct:doSomethingElse")                           // 3 (POJO is still there)
    

    请注意,此自定义处理器与标准 WireTap() 的同步模拟并不完全一致,后者完全是 in-only,因为此处理器重新抛出目标路由上发生的任何未处理异常 - 但消息本身保持不变.这是我的要求,因为我想做的是在另一条路由上同步执行一些其他处理,并在失败时收到通知,否则我的主路由上的消息不会受到影响(相当于调用 void 方法在程序代码中)。

    【讨论】:

      【解决方案5】:

      您可以使用窃听中的 copy="true" 选项来复制 http://camel.apache.org/wire-tap.html 中提到的交换,或者您可以创建自己的处理器来做同样的事情。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2023-03-06
        • 1970-01-01
        • 2015-08-01
        • 2017-05-19
        • 2017-04-18
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多