【问题标题】:Integration pattern : how to sync processing message received from multiple systems集成模式:如何同步从多个系统接收到的处理消息
【发布时间】:2018-02-25 12:40:14
【问题描述】:

我正在构建一个系统,该系统将通过消息代理(当前为 JMS)接收来自不同系统的消息。来自所有发送者系统的所有消息都有一个 deviceId,并且消息的接收没有顺序。 例如,系统 A 可以发送 deviceId=1 的消息,系统 b 可以发送 deviceId=2 的消息。

我的目标是不开始处理有关同一 deviceId 的消息,除非我从所有具有相同 deviceId 的发件人那里收到所有消息。

例如,如果我有 3 个系统 A、B 和 C 向我的系统发送消息:

System A sends messageA1 with deviceId=1
System B sends messageB1 with deviceId=1
System C sends messageC1 with deviceId=3
System C sends messageC2 with deviceId=1 <--- here I should start processing of messageA1, messageB1 and messageC2 because they are having the same deviceID 1.

是否应该通过在我的系统中使用某种同步机制、消息代理或 spring-integration/apache camel 等集成框架来解决此问题?

【问题讨论】:

  • 服务是否从其中一个系统接收到几条具有相同 deviceId 的消息?您想一一处理单独的消息吗?示例(deviceId=1):messageA1, messageA1, messageA1, messageC1, messageC1, messageB1 -> 进程已启动
  • 不,没关系......重要的是当所有3条消息都到达时触发处理,因为它们包含我必须处理的数据

标签: java design-patterns apache-camel spring-integration integration-patterns


【解决方案1】:

与聚合器类似的解决方案(@Artem Bilan 提到的)也可以在 Camel 中实现,使用自定义 AggregationStrategy 并使用 Exchange.AGGREGATION_COMPLETE_CURRENT_GROUP 属性控制聚合器完成。

以下可能是一个很好的起点。 (You can find the sample project with tests here)

路线:

from("direct:start")
    .log(LoggingLevel.INFO, "Received ${headers.system}${headers.deviceId}")
    .aggregate(header("deviceId"), new SignalAggregationStrategy(3))
    .log(LoggingLevel.INFO, "Signaled body: ${body}")
    .to("direct:result");

SignalAggregationStrategy.java

public class SignalAggregationStrategy extends GroupedExchangeAggregationStrategy implements Predicate {

    private int numberOfSystems;

    public SignalAggregationStrategy(int numberOfSystems) {
        this.numberOfSystems = numberOfSystems;
    }

    @Override
    public Exchange aggregate(Exchange oldExchange, Exchange newExchange) {
        Exchange exchange = super.aggregate(oldExchange, newExchange);

        List<Exchange> aggregatedExchanges = exchange.getProperty("CamelGroupedExchange", List.class);

        // Complete aggregation if we have "numberOfSystems" (currently 3) different messages (where "system" headers are different)
        // https://github.com/apache/camel/blob/master/camel-core/src/main/docs/eips/aggregate-eip.adoc#completing-current-group-decided-from-the-aggregationstrategy
        if (numberOfSystems == aggregatedExchanges.stream().map(e -> e.getIn().getHeader("system", String.class)).distinct().count()) {
            exchange.setProperty(Exchange.AGGREGATION_COMPLETE_CURRENT_GROUP, true);
        }

        return exchange;
    }

    @Override
    public boolean matches(Exchange exchange) {
        // make it infinite (4th bullet point @ https://github.com/apache/camel/blob/master/camel-core/src/main/docs/eips/aggregate-eip.adoc#about-completion)
        return false;
    }
}

希望对你有帮助!

【讨论】:

    【解决方案2】:

    您可以在 Apache Camel 中使用缓存组件执行此操作。我认为有 EHCache 组件。

    基本上:

    1. 您收到一条带有给定 deviceId 的消息,例如 deviceId1。
    2. 您在缓存中查找已收到 deviceId1 的哪些消息。
    3. 只要您没有收到所有三个,您就可以将当前系统/消息添加到缓存中。
    4. 一旦所有消息都在那里,您就可以处理并清除缓存。

    然后,您当然可以将每条传入消息路由到特定的基于 deviceId 的队列以进行临时存储。这可以是 JMS、ActiveMQ 或类似的东西。

    【讨论】:

      【解决方案3】:

      Spring Integration 为此类任务提供了组件 - 在收集到整个组之前不要发出。它的名字是Aggregator。你的deviceId 绝对是correlationKey。 releaseStrategy 实际上可能取决于系统的数量 - 在继续下一步之前您正在等待多少 deviceId1 消息。

      【讨论】:

        猜你喜欢
        • 2014-10-27
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2011-06-18
        • 2018-01-19
        • 2021-02-08
        • 1970-01-01
        • 2015-06-22
        相关资源
        最近更新 更多