关于correlationId、sequenceNumber及聚合排序的技术问询
Spring Integration Scatter-Gather 聚合相关问题
我正在使用包含三个并行流的Scatter-Gather模式,发现聚合后的groupMessage按响应到达顺序排列。现咨询:
- correlationId和sequenceNumber是什么,在聚合中扮演什么角色?
- 是否可以通过这两个字段修改aggregate()方法中分组消息的索引?
此外,运行时group为列表,我无法确定索引,需将三个响应作为参数做进一步处理,相关代码及日志如下:
代码片段
public IntegrationFlow flow() { return flow -> flow.split() .channel(c -> c.executor(Executors.newCachedThreadPool())) .scatterGather( scatterer -> scatterer .applySequence(true) .recipientFlow(flow1()) .recipientFlow(flow2()) .recipientFlow(flow3()), gatherer -> gatherer.releaseLockBeforeSend(true)) .log("Gatherer") .aggregate(aggregatingAndFormingRequestForFlow4()) .to(flow4()); } public MessageGroupProcessor aggregatingAndFormingRequestForFlow4() { return group -> { System.out.println("Get Messages: " + group.getMessages()); // 这里需要将三个响应作为参数做后续处理,但运行时group是列表,无法确定索引 prepareRequest(); }; }
日志信息
INFO [pool-12-thread-1]Gatherer:GenericMessage [payload=[response from flow1, response from flow2, response from flow3 ], headers={sequenceNumber=1, sequenceDetails=[[f4da9a6f-f27d-91ac-cbef-5106eb614d62, 1, 1]], errorChannel=org.springframework.messaging.core.GenericMessagingTemplate$TemporaryReplyChannel@1b6efbf0, sequenceSize=1, replyChannel=org.springframework.messaging.core.GenericMessagingTemplate$TemporaryReplyChannel@1b6efbf0,correlationId=f4da9a6f-f27d-91ac-cbef-5106eb614d62, id=378f5406-ce9e-19fb-c909-5d5c7c4093c0, timestamp=1662660858032}] Get Messages: [payload=[response from flow1, response from flow2, response from flow3 ], headers={sequenceNumber=1, sequenceDetails=[[f4da9a6f-f27d-91ac-cbef-5106eb614d62, 1, 1]], errorChannel=org.springframework.messaging.core.GenericMessagingTemplate$TemporaryReplyChannel@1b6efbf0, sequenceSize=1, replyChannel=org.springframework.messaging.core.GenericMessagingTemplate$TemporaryReplyChannel@1b6efbf0,correlationId=f4da9a6f-f27d-91ac-cbef-5106eb614d62, id=378f5406-ce9e-19fb-c909-5d5c7c4093c0, timestamp=1662660858032}]
可见日志中是三个响应组成的数组。
问题解答
1. correlationId 和 sequenceNumber 的定义与聚合中的角色
- correlationId:消息分组的唯一标识,用于将分散出去的请求消息及其响应关联到同一个聚合组。同一Scatter-Gather请求下的所有消息(包括分散的请求和对应响应)都会携带相同的correlationId,聚合器通过这个ID识别哪些消息属于同一组,避免跨请求混淆。
- sequenceNumber:标识当前消息在所属分组内的序号。开启
applySequence(true)后,Scatterer会为每个发送到下游流的请求分配递增的sequenceNumber(比如三个recipientFlow对应1、2、3),响应消息会继承该序号。你的日志中sequenceNumber显示为1,是因为gatherer已经将三个响应聚合为单个消息,这个聚合后的消息作为整体的序号是1。
2. 是否可以通过这两个字段修改aggregate()中分组消息的索引?
不能直接修改索引,但可以利用这两个字段实现消息排序。默认聚合器按消息到达顺序排列,你可以在aggregate()或scatterGather()的gatherer配置中,基于原始响应消息的sequenceNumber自定义排序逻辑,让列表顺序与recipientFlow的定义顺序一致。
解决无法确定索引的问题
从日志可以看到,group.getMessages()返回的是gatherer已聚合好的单个消息(payload为三个响应的数组),无需处理分组索引,直接提取数组元素即可:
修改aggregatingAndFormingRequestForFlow4()方法:
public MessageGroupProcessor aggregatingAndFormingRequestForFlow4() { return group -> { // 获取gatherer聚合后的消息 Message<?> aggregatedMessage = group.getMessages().iterator().next(); // 提取payload中的三个响应 List<?> responses = (List<?>) aggregatedMessage.getPayload(); Object flow1Response = responses.get(0); Object flow2Response = responses.get(1); Object flow3Response = responses.get(2); // 用三个响应做后续处理 prepareRequest(flow1Response, flow2Response, flow3Response); // 返回处理后的结果消息 return MessageBuilder.withPayload(/* 处理后的结果 */).build(); }; }
如果要确保gatherer聚合后的数组顺序与recipientFlow定义顺序一致(而非响应到达顺序),可配置gatherer的排序策略:
.scatterGather( scatterer -> scatterer .applySequence(true) .recipientFlow(flow1()) .recipientFlow(flow2()) .recipientFlow(flow3()), gatherer -> gatherer.releaseLockBeforeSend(true) .order(group -> group.getMessages() .stream() .sorted(Comparator.comparingInt(msg -> msg.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER, Integer.class))) .collect(Collectors.toList())) )
这样gatherer会根据响应消息的sequenceNumber排序,确保数组顺序与三个流的定义顺序一致,之后就能安全按索引取对应流的响应。
内容的提问来源于stack exchange,提问作者Somnath Mukherjee
相关产品推荐
相关产品推荐

