【发布时间】: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