【问题标题】:How to use Spring Integration Aggregator for 'look-ahead' releasing?如何使用 Spring Integration Aggregator 进行“前瞻”发布?
【发布时间】:2017-03-02 17:29:06
【问题描述】:

问题

我正在使用 Spring Integration 开发应用程序(版本 5.0.0.M2,因为我需要动态流注册功能)并且在某些时候我必须根据它们的 correlationId 标头聚合消息。标头标记了消息之间的一些连续子序列,即它仅在新子序列开始时发生变化。聚合器的目的是将子序列转换为单独的“分组”消息。
消息的来源是外部且不可预测的,因此它无法为发出的消息提供sequenceSize 标头。因此,直到下一条消息带有另一个correlationId 标头(或者timeout ms 没有下一条消息,这意味着输入结束),才能决定是否释放当前的累积组。部分释放是不可接受的。

问题是 Spring Integration ReleaseStrategy 旨在仅发布 当前 累积组(包括当前消息),而我的任务需要“向前看”下一条消息以了解当前组是否完整。从下一条消息的角度来看,同样的问题可能被认为是“回顾”发布。

所以问题是:有没有办法根据下一条消息的标头(不包括后者)从聚合器中释放消息组?

尝试过的解决方案

我研究了org.springframework.integration.aggregator.AbstractCorrelatingMessageHandler 的源代码,目前我发现的唯一解决方法是一个相当肮脏的技巧——我继承自AggregatingMessageHandler 并覆盖了handleMessageInternal 方法。在该方法中,我将当前消息的correlationId 与上一次调用中保存的消息进行比较。如果它不同,我会从商店中提取上一个组并用它调用forceComplete。然后(以任何方式)我将当前消息处理委托给父级的handleMessageInternal 方法。这是sn-p的代码:

  @Override
  protected void handleMessageInternal(Message<?> message) throws Exception {
    Long currentGroupId = message.getHeaders().get(CORRELATION_ID, Long.class);

    boolean needToReleasePreviousGroup = ((previousGroupId != null) && !previousGroupId.equals(currentGroupId));
    if (needToReleasePreviousGroup) {
      MessageGroup previousGroup = getMessageStore().getMessageGroup(previousGroupId);
      if (previousGroup.size() != 0) {
        forceComplete(previousGroup);
      } else {
        log.debug("Previous group with id={} has been already released. Skip.", previousGroupId);
      }
    }

    super.handleMessageInternal(message);
    previousGroupId = currentGroupId;   // do unconditionally as we'll check that group on the next step anyway
  }

我意识到这个解决方案使聚合器本身有状态(虽然它不应该是)并且不适用于许多其他场景。此外,它使应用的代码与框架的内部耦合。

如果有人指出更好的解决方案,我会很高兴。我愿意提供有关问题或我使用的解决方案的更多详细信息。

【问题讨论】:

    标签: java spring spring-integration


    【解决方案1】:

    期待您的previousGroupId 解决方案和您的描述为:

    “向前看”查看下一条消息,以确定当前组是否完整。

    “当前”这个词是这里的关键。

    所以,对我来说,这意味着您的聚合器是 single group,您必须释放当前状态并根据当前消息的状态开始一个新状态。

    为此,我将聚合器设为真正的单组 - correlationStrategy 应该返回一些常量,例如1.

    是否释放当前组确实是ReleaseStrategy 的责任。最新消息将是MessageGroup.getMessages() 中的最后一条消息。因此,您收到该消息,决定是否需要发布并继续使用自定义 MessageGroupProcessor 以从聚合器生成累积结果。

    在那MessageGroupProcessor你删除最新的“有罪”消息不要将其包含在结果中并将其发送回聚合器以形成一个新的消息组。

    这对你有意义吗?

    【讨论】:

    • 是的,这似乎是我需要的。但我仍然看不到如何'将其发送回聚合器'?直接从MessageGroupProcessor 执行此操作是毫无意义的,因为当前组尚未发布。在我看来,通过一些异步重新提交到聚合器的输入通道来完成它似乎增加了应用程序的复杂性并产生了过多的工作。是否有更简洁的方式将消息重新提交给聚合器(或任何其他消息处理程序)?
    • 确实如此。您必须将该消息发送到QueueChannel 以进行线程转移并让当前线程释放组并清理资源。在correlationKey 的锁未被释放之前,该队列上的轮询器将无法将该消息放置到聚合器。
    • 好的,我会尝试应用这种方法。锁定correlationKey 是我在想象它如何工作时没有考虑到的难题的一部分。得到结果后我会立即分享。
    • 注意你可以将LockRegistry注入聚合器,在其他地方也可以使用相同的lock。
    • 我已经 added a comment 提出了我在这个问题中遇到的类似情况的建议。
    【解决方案2】:

    对于那些试图解决相同或类似问题的人,有一个 Github gist 包含基于公认答案构建的聚合器配置的源代码。

    特意为问题简化了源代码。真正的还可以通过超时或组大小阈值来释放组。如有需要,也准备分享。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-10-02
      • 2016-04-04
      • 1970-01-01
      • 1970-01-01
      • 2015-01-16
      相关资源
      最近更新 更多