Fix Sonar complexity issues
Address a few complexity issues.
This commit is contained in:
committed by
Artem Bilan
parent
9072e925e9
commit
c570bee197
@@ -26,6 +26,7 @@ import java.util.UUID;
|
||||
import org.springframework.amqp.core.MessageDeliveryMode;
|
||||
import org.springframework.amqp.core.MessageProperties;
|
||||
import org.springframework.amqp.support.AmqpHeaders;
|
||||
import org.springframework.amqp.utils.JavaUtils;
|
||||
import org.springframework.integration.IntegrationMessageHeaderAccessor;
|
||||
import org.springframework.integration.mapping.AbstractHeaderMapper;
|
||||
import org.springframework.integration.mapping.support.JsonHeaders;
|
||||
@@ -107,86 +108,55 @@ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper<MessagePropert
|
||||
protected Map<String, Object> extractStandardHeaders(MessageProperties amqpMessageProperties) {
|
||||
Map<String, Object> headers = new HashMap<String, Object>();
|
||||
try {
|
||||
String appId = amqpMessageProperties.getAppId();
|
||||
if (StringUtils.hasText(appId)) {
|
||||
headers.put(AmqpHeaders.APP_ID, appId);
|
||||
}
|
||||
String clusterId = amqpMessageProperties.getClusterId();
|
||||
if (StringUtils.hasText(clusterId)) {
|
||||
headers.put(AmqpHeaders.CLUSTER_ID, clusterId);
|
||||
}
|
||||
String contentEncoding = amqpMessageProperties.getContentEncoding();
|
||||
if (StringUtils.hasText(contentEncoding)) {
|
||||
headers.put(AmqpHeaders.CONTENT_ENCODING, contentEncoding);
|
||||
}
|
||||
JavaUtils.INSTANCE
|
||||
.acceptIfNotNull(AmqpHeaders.APP_ID, amqpMessageProperties.getAppId(),
|
||||
(key, value) -> headers.put(key, value))
|
||||
.acceptIfNotNull(AmqpHeaders.CLUSTER_ID, amqpMessageProperties.getClusterId(),
|
||||
(key, value) -> headers.put(key, value))
|
||||
.acceptIfNotNull(AmqpHeaders.CONTENT_ENCODING, amqpMessageProperties.getContentEncoding(),
|
||||
(key, value) -> headers.put(key, value));
|
||||
long contentLength = amqpMessageProperties.getContentLength();
|
||||
if (contentLength > 0) {
|
||||
headers.put(AmqpHeaders.CONTENT_LENGTH, contentLength);
|
||||
}
|
||||
String contentType = amqpMessageProperties.getContentType();
|
||||
if (StringUtils.hasText(contentType)) {
|
||||
headers.put(AmqpHeaders.CONTENT_TYPE, contentType);
|
||||
}
|
||||
String correlationId = amqpMessageProperties.getCorrelationId();
|
||||
if (StringUtils.hasText(correlationId)) {
|
||||
headers.put(AmqpHeaders.CORRELATION_ID, correlationId);
|
||||
}
|
||||
MessageDeliveryMode receivedDeliveryMode = amqpMessageProperties.getReceivedDeliveryMode();
|
||||
if (receivedDeliveryMode != null) {
|
||||
headers.put(AmqpHeaders.RECEIVED_DELIVERY_MODE, receivedDeliveryMode);
|
||||
}
|
||||
JavaUtils.INSTANCE
|
||||
.acceptIfCondition(contentLength > 0, AmqpHeaders.CONTENT_LENGTH, contentLength,
|
||||
(key, value) -> headers.put(key, value))
|
||||
.acceptIfHasText(AmqpHeaders.CONTENT_TYPE, amqpMessageProperties.getContentType(),
|
||||
(key, value) -> headers.put(key, value))
|
||||
.acceptIfHasText(AmqpHeaders.CORRELATION_ID, amqpMessageProperties.getCorrelationId(),
|
||||
(key, value) -> headers.put(key, value))
|
||||
.acceptIfNotNull(AmqpHeaders.RECEIVED_DELIVERY_MODE, amqpMessageProperties.getReceivedDeliveryMode(),
|
||||
(key, value) -> headers.put(key, value));
|
||||
long deliveryTag = amqpMessageProperties.getDeliveryTag();
|
||||
if (deliveryTag > 0) {
|
||||
headers.put(AmqpHeaders.DELIVERY_TAG, deliveryTag);
|
||||
}
|
||||
String expiration = amqpMessageProperties.getExpiration();
|
||||
if (StringUtils.hasText(expiration)) {
|
||||
headers.put(AmqpHeaders.EXPIRATION, expiration);
|
||||
}
|
||||
JavaUtils.INSTANCE
|
||||
.acceptIfCondition(deliveryTag > 0, AmqpHeaders.DELIVERY_TAG, deliveryTag,
|
||||
(key, value) -> headers.put(key, value))
|
||||
.acceptIfHasText(AmqpHeaders.EXPIRATION, amqpMessageProperties.getExpiration(),
|
||||
(key, value) -> headers.put(key, value));
|
||||
Integer messageCount = amqpMessageProperties.getMessageCount();
|
||||
if (messageCount != null && messageCount > 0) {
|
||||
headers.put(AmqpHeaders.MESSAGE_COUNT, messageCount);
|
||||
}
|
||||
String messageId = amqpMessageProperties.getMessageId();
|
||||
if (StringUtils.hasText(messageId)) {
|
||||
headers.put(AmqpHeaders.MESSAGE_ID, messageId);
|
||||
}
|
||||
JavaUtils.INSTANCE
|
||||
.acceptIfCondition(messageCount != null && messageCount > 0, AmqpHeaders.MESSAGE_COUNT, messageCount,
|
||||
(key, value) -> headers.put(key, value))
|
||||
.acceptIfHasText(AmqpHeaders.MESSAGE_ID, amqpMessageProperties.getMessageId(),
|
||||
(key, value) -> headers.put(key, value));
|
||||
Integer priority = amqpMessageProperties.getPriority();
|
||||
if (priority != null && priority > 0) {
|
||||
headers.put(IntegrationMessageHeaderAccessor.PRIORITY, priority);
|
||||
}
|
||||
Integer receivedDelay = amqpMessageProperties.getReceivedDelay();
|
||||
if (receivedDelay != null) {
|
||||
headers.put(AmqpHeaders.RECEIVED_DELAY, receivedDelay);
|
||||
}
|
||||
String receivedExchange = amqpMessageProperties.getReceivedExchange();
|
||||
if (StringUtils.hasText(receivedExchange)) {
|
||||
headers.put(AmqpHeaders.RECEIVED_EXCHANGE, receivedExchange);
|
||||
}
|
||||
String receivedRoutingKey = amqpMessageProperties.getReceivedRoutingKey();
|
||||
if (StringUtils.hasText(receivedRoutingKey)) {
|
||||
headers.put(AmqpHeaders.RECEIVED_ROUTING_KEY, receivedRoutingKey);
|
||||
}
|
||||
Boolean redelivered = amqpMessageProperties.isRedelivered();
|
||||
if (redelivered != null) {
|
||||
headers.put(AmqpHeaders.REDELIVERED, redelivered);
|
||||
}
|
||||
String replyTo = amqpMessageProperties.getReplyTo();
|
||||
if (replyTo != null) {
|
||||
headers.put(AmqpHeaders.REPLY_TO, replyTo);
|
||||
}
|
||||
Date timestamp = amqpMessageProperties.getTimestamp();
|
||||
if (timestamp != null) {
|
||||
headers.put(AmqpHeaders.TIMESTAMP, timestamp);
|
||||
}
|
||||
String type = amqpMessageProperties.getType();
|
||||
if (StringUtils.hasText(type)) {
|
||||
headers.put(AmqpHeaders.TYPE, type);
|
||||
}
|
||||
String userId = amqpMessageProperties.getReceivedUserId();
|
||||
if (StringUtils.hasText(userId)) {
|
||||
headers.put(AmqpHeaders.RECEIVED_USER_ID, userId);
|
||||
}
|
||||
JavaUtils.INSTANCE
|
||||
.acceptIfCondition(priority != null && priority > 0, IntegrationMessageHeaderAccessor.PRIORITY,
|
||||
priority, (key, value) -> headers.put(key, value))
|
||||
.acceptIfNotNull(AmqpHeaders.RECEIVED_DELAY, amqpMessageProperties.getReceivedDelay(),
|
||||
(key, value) -> headers.put(key, value))
|
||||
.acceptIfNotNull(AmqpHeaders.RECEIVED_EXCHANGE, amqpMessageProperties.getReceivedExchange(),
|
||||
(key, value) -> headers.put(key, value))
|
||||
.acceptIfHasText(AmqpHeaders.RECEIVED_ROUTING_KEY, amqpMessageProperties.getReceivedRoutingKey(),
|
||||
(key, value) -> headers.put(key, value))
|
||||
.acceptIfNotNull(AmqpHeaders.REDELIVERED, amqpMessageProperties.isRedelivered(),
|
||||
(key, value) -> headers.put(key, value))
|
||||
.acceptIfNotNull(AmqpHeaders.REPLY_TO, amqpMessageProperties.getReplyTo(),
|
||||
(key, value) -> headers.put(key, value))
|
||||
.acceptIfNotNull(AmqpHeaders.TIMESTAMP, amqpMessageProperties.getTimestamp(),
|
||||
(key, value) -> headers.put(key, value))
|
||||
.acceptIfHasText(AmqpHeaders.TYPE, amqpMessageProperties.getType(),
|
||||
(key, value) -> headers.put(key, value))
|
||||
.acceptIfHasText(AmqpHeaders.RECEIVED_USER_ID, amqpMessageProperties.getReceivedUserId(),
|
||||
(key, value) -> headers.put(key, value));
|
||||
|
||||
for (String jsonHeader : JsonHeaders.HEADERS) {
|
||||
Object value = amqpMessageProperties.getHeaders().get(jsonHeader.replaceFirst(JsonHeaders.PREFIX, ""));
|
||||
@@ -228,54 +198,30 @@ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper<MessagePropert
|
||||
@Override
|
||||
protected void populateStandardHeaders(@Nullable Map<String, Object> allHeaders, Map<String, Object> headers,
|
||||
MessageProperties amqpMessageProperties) {
|
||||
String appId = getHeaderIfAvailable(headers, AmqpHeaders.APP_ID, String.class);
|
||||
if (StringUtils.hasText(appId)) {
|
||||
amqpMessageProperties.setAppId(appId);
|
||||
}
|
||||
String clusterId = getHeaderIfAvailable(headers, AmqpHeaders.CLUSTER_ID, String.class);
|
||||
if (StringUtils.hasText(clusterId)) {
|
||||
amqpMessageProperties.setClusterId(clusterId);
|
||||
}
|
||||
String contentEncoding = getHeaderIfAvailable(headers, AmqpHeaders.CONTENT_ENCODING, String.class);
|
||||
if (StringUtils.hasText(contentEncoding)) {
|
||||
amqpMessageProperties.setContentEncoding(contentEncoding);
|
||||
}
|
||||
Long contentLength = getHeaderIfAvailable(headers, AmqpHeaders.CONTENT_LENGTH, Long.class);
|
||||
if (contentLength != null) {
|
||||
amqpMessageProperties.setContentLength(contentLength);
|
||||
}
|
||||
String contentType = this.extractContentTypeAsString(headers);
|
||||
|
||||
if (StringUtils.hasText(contentType)) {
|
||||
amqpMessageProperties.setContentType(contentType);
|
||||
}
|
||||
|
||||
String correlationId = getHeaderIfAvailable(headers, AmqpHeaders.CORRELATION_ID, String.class);
|
||||
if (StringUtils.hasText(correlationId)) {
|
||||
amqpMessageProperties.setCorrelationId(correlationId);
|
||||
}
|
||||
|
||||
Integer delay = getHeaderIfAvailable(headers, AmqpHeaders.DELAY, Integer.class);
|
||||
if (delay != null) {
|
||||
amqpMessageProperties.setDelay(delay);
|
||||
}
|
||||
MessageDeliveryMode deliveryMode = getHeaderIfAvailable(headers, AmqpHeaders.DELIVERY_MODE,
|
||||
MessageDeliveryMode.class);
|
||||
if (deliveryMode != null) {
|
||||
amqpMessageProperties.setDeliveryMode(deliveryMode);
|
||||
}
|
||||
Long deliveryTag = getHeaderIfAvailable(headers, AmqpHeaders.DELIVERY_TAG, Long.class);
|
||||
if (deliveryTag != null) {
|
||||
amqpMessageProperties.setDeliveryTag(deliveryTag);
|
||||
}
|
||||
String expiration = getHeaderIfAvailable(headers, AmqpHeaders.EXPIRATION, String.class);
|
||||
if (StringUtils.hasText(expiration)) {
|
||||
amqpMessageProperties.setExpiration(expiration);
|
||||
}
|
||||
Integer messageCount = getHeaderIfAvailable(headers, AmqpHeaders.MESSAGE_COUNT, Integer.class);
|
||||
if (messageCount != null) {
|
||||
amqpMessageProperties.setMessageCount(messageCount);
|
||||
}
|
||||
JavaUtils.INSTANCE
|
||||
.acceptIfHasText(getHeaderIfAvailable(headers, AmqpHeaders.APP_ID, String.class),
|
||||
appId -> amqpMessageProperties.setAppId(appId))
|
||||
.acceptIfHasText(getHeaderIfAvailable(headers, AmqpHeaders.CLUSTER_ID, String.class),
|
||||
clusterId -> amqpMessageProperties.setClusterId(clusterId))
|
||||
.acceptIfHasText(getHeaderIfAvailable(headers, AmqpHeaders.CONTENT_ENCODING, String.class),
|
||||
contentEncoding -> amqpMessageProperties.setContentEncoding(contentEncoding))
|
||||
.acceptIfNotNull(getHeaderIfAvailable(headers, AmqpHeaders.CONTENT_LENGTH, Long.class),
|
||||
contentLength -> amqpMessageProperties.setContentLength(contentLength))
|
||||
.acceptIfHasText(this.extractContentTypeAsString(headers),
|
||||
contentType -> amqpMessageProperties.setContentType(contentType))
|
||||
.acceptIfHasText(getHeaderIfAvailable(headers, AmqpHeaders.CORRELATION_ID, String.class),
|
||||
correlationId -> amqpMessageProperties.setCorrelationId(correlationId))
|
||||
.acceptIfNotNull(getHeaderIfAvailable(headers, AmqpHeaders.DELAY, Integer.class),
|
||||
delay -> amqpMessageProperties.setDelay(delay))
|
||||
.acceptIfNotNull(getHeaderIfAvailable(headers, AmqpHeaders.DELIVERY_MODE, MessageDeliveryMode.class),
|
||||
deliveryMode -> amqpMessageProperties.setDeliveryMode(deliveryMode))
|
||||
.acceptIfNotNull(getHeaderIfAvailable(headers, AmqpHeaders.DELIVERY_TAG, Long.class),
|
||||
deliveryTag -> amqpMessageProperties.setDeliveryTag(deliveryTag))
|
||||
.acceptIfHasText(getHeaderIfAvailable(headers, AmqpHeaders.EXPIRATION, String.class),
|
||||
expiration -> amqpMessageProperties.setExpiration(expiration))
|
||||
.acceptIfNotNull(getHeaderIfAvailable(headers, AmqpHeaders.MESSAGE_COUNT, Integer.class),
|
||||
messageCount -> amqpMessageProperties.setMessageCount(messageCount));
|
||||
String messageId = getHeaderIfAvailable(headers, AmqpHeaders.MESSAGE_ID, String.class);
|
||||
if (StringUtils.hasText(messageId)) {
|
||||
amqpMessageProperties.setMessageId(messageId);
|
||||
@@ -286,26 +232,17 @@ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper<MessagePropert
|
||||
amqpMessageProperties.setMessageId(id.toString());
|
||||
}
|
||||
}
|
||||
Integer priority = getHeaderIfAvailable(headers, IntegrationMessageHeaderAccessor.PRIORITY, Integer.class);
|
||||
if (priority != null) {
|
||||
amqpMessageProperties.setPriority(priority);
|
||||
}
|
||||
String receivedExchange = getHeaderIfAvailable(headers, AmqpHeaders.RECEIVED_EXCHANGE, String.class);
|
||||
if (StringUtils.hasText(receivedExchange)) {
|
||||
amqpMessageProperties.setReceivedExchange(receivedExchange);
|
||||
}
|
||||
String receivedRoutingKey = getHeaderIfAvailable(headers, AmqpHeaders.RECEIVED_ROUTING_KEY, String.class);
|
||||
if (StringUtils.hasText(receivedRoutingKey)) {
|
||||
amqpMessageProperties.setReceivedRoutingKey(receivedRoutingKey);
|
||||
}
|
||||
Boolean redelivered = getHeaderIfAvailable(headers, AmqpHeaders.REDELIVERED, Boolean.class);
|
||||
if (redelivered != null) {
|
||||
amqpMessageProperties.setRedelivered(redelivered);
|
||||
}
|
||||
String replyTo = getHeaderIfAvailable(headers, AmqpHeaders.REPLY_TO, String.class);
|
||||
if (replyTo != null) {
|
||||
amqpMessageProperties.setReplyTo(replyTo);
|
||||
}
|
||||
JavaUtils.INSTANCE
|
||||
.acceptIfNotNull(getHeaderIfAvailable(headers, IntegrationMessageHeaderAccessor.PRIORITY, Integer.class),
|
||||
priority -> amqpMessageProperties.setPriority(priority))
|
||||
.acceptIfHasText(getHeaderIfAvailable(headers, AmqpHeaders.RECEIVED_EXCHANGE, String.class),
|
||||
receivedExchange -> amqpMessageProperties.setReceivedExchange(receivedExchange))
|
||||
.acceptIfHasText(getHeaderIfAvailable(headers, AmqpHeaders.RECEIVED_ROUTING_KEY, String.class),
|
||||
receivedRoutingKey -> amqpMessageProperties.setReceivedRoutingKey(receivedRoutingKey))
|
||||
.acceptIfNotNull(getHeaderIfAvailable(headers, AmqpHeaders.REDELIVERED, Boolean.class),
|
||||
redelivered -> amqpMessageProperties.setRedelivered(redelivered))
|
||||
.acceptIfNotNull(getHeaderIfAvailable(headers, AmqpHeaders.REPLY_TO, String.class),
|
||||
replyTo -> amqpMessageProperties.setReplyTo(replyTo));
|
||||
Date timestamp = getHeaderIfAvailable(headers, AmqpHeaders.TIMESTAMP, Date.class);
|
||||
if (timestamp != null) {
|
||||
amqpMessageProperties.setTimestamp(timestamp);
|
||||
@@ -316,14 +253,11 @@ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper<MessagePropert
|
||||
amqpMessageProperties.setTimestamp(new Date(ts));
|
||||
}
|
||||
}
|
||||
String type = getHeaderIfAvailable(headers, AmqpHeaders.TYPE, String.class);
|
||||
if (type != null) {
|
||||
amqpMessageProperties.setType(type);
|
||||
}
|
||||
String userId = getHeaderIfAvailable(headers, AmqpHeaders.USER_ID, String.class);
|
||||
if (StringUtils.hasText(userId)) {
|
||||
amqpMessageProperties.setUserId(userId);
|
||||
}
|
||||
JavaUtils.INSTANCE
|
||||
.acceptIfNotNull(getHeaderIfAvailable(headers, AmqpHeaders.TYPE, String.class),
|
||||
type -> amqpMessageProperties.setType(type))
|
||||
.acceptIfNotNull(getHeaderIfAvailable(headers, AmqpHeaders.USER_ID, String.class),
|
||||
userId -> amqpMessageProperties.setUserId(userId));
|
||||
|
||||
Map<String, String> jsonHeaders = new HashMap<String, String>();
|
||||
|
||||
@@ -346,14 +280,11 @@ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper<MessagePropert
|
||||
amqpMessageProperties.getHeaders().putAll(jsonHeaders);
|
||||
}
|
||||
|
||||
String replyCorrelation = getHeaderIfAvailable(headers, AmqpHeaders.SPRING_REPLY_CORRELATION, String.class);
|
||||
if (StringUtils.hasLength(replyCorrelation)) {
|
||||
amqpMessageProperties.setHeader("spring_reply_correlation", replyCorrelation);
|
||||
}
|
||||
String replyToStack = getHeaderIfAvailable(headers, AmqpHeaders.SPRING_REPLY_TO_STACK, String.class);
|
||||
if (StringUtils.hasLength(replyToStack)) {
|
||||
amqpMessageProperties.setHeader("spring_reply_to", replyToStack);
|
||||
}
|
||||
JavaUtils.INSTANCE
|
||||
.acceptIfHasText(getHeaderIfAvailable(headers, AmqpHeaders.SPRING_REPLY_CORRELATION, String.class),
|
||||
replyCorrelation -> amqpMessageProperties.setHeader("spring_reply_correlation", replyCorrelation))
|
||||
.acceptIfHasText(getHeaderIfAvailable(headers, AmqpHeaders.SPRING_REPLY_TO_STACK, String.class),
|
||||
replyToStack -> amqpMessageProperties.setHeader("spring_reply_to", replyToStack));
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -105,6 +105,18 @@ public abstract class AbstractAggregatingMessageGroupProcessor implements Messag
|
||||
*/
|
||||
protected Map<String, Object> aggregateHeaders(MessageGroup group) {
|
||||
Map<String, Object> aggregatedHeaders = new HashMap<>();
|
||||
Set<String> conflictKeys = doAggregateHeaders(group, aggregatedHeaders);
|
||||
for (String keyToRemove : conflictKeys) {
|
||||
if (this.logger.isDebugEnabled()) {
|
||||
this.logger.debug("Excluding header '" + keyToRemove + "' upon aggregation due to conflict(s) "
|
||||
+ "in MessageGroup with correlation key: " + group.getGroupId());
|
||||
}
|
||||
aggregatedHeaders.remove(keyToRemove);
|
||||
}
|
||||
return aggregatedHeaders;
|
||||
}
|
||||
|
||||
private Set<String> doAggregateHeaders(MessageGroup group, Map<String, Object> aggregatedHeaders) {
|
||||
Set<String> conflictKeys = new HashSet<>();
|
||||
for (Message<?> message : group.getMessages()) {
|
||||
for (Entry<String, Object> entry : message.getHeaders().entrySet()) {
|
||||
@@ -126,14 +138,7 @@ public abstract class AbstractAggregatingMessageGroupProcessor implements Messag
|
||||
}
|
||||
}
|
||||
}
|
||||
for (String keyToRemove : conflictKeys) {
|
||||
if (this.logger.isDebugEnabled()) {
|
||||
this.logger.debug("Excluding header '" + keyToRemove + "' upon aggregation due to conflict(s) "
|
||||
+ "in MessageGroup with correlation key: " + group.getGroupId());
|
||||
}
|
||||
aggregatedHeaders.remove(keyToRemove);
|
||||
}
|
||||
return aggregatedHeaders;
|
||||
return conflictKeys;
|
||||
}
|
||||
|
||||
protected abstract Object aggregatePayloads(MessageGroup group, Map<String, Object> defaultHeaders);
|
||||
|
||||
@@ -643,7 +643,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP
|
||||
afterRelease(group, completedMessages);
|
||||
}
|
||||
|
||||
protected void forceComplete(MessageGroup group) {
|
||||
protected void forceComplete(MessageGroup group) { // NOSONAR Complexity
|
||||
Object correlationKey = group.getGroupId();
|
||||
// UUIDConverter is no-op if already converted
|
||||
UUID groupId = UUIDConverter.getUUID(correlationKey);
|
||||
@@ -737,7 +737,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP
|
||||
}
|
||||
}
|
||||
}
|
||||
catch (InterruptedException ie) {
|
||||
catch (@SuppressWarnings("unused") InterruptedException ie) {
|
||||
Thread.currentThread().interrupt();
|
||||
this.logger.debug("Thread was interrupted while trying to obtain lock");
|
||||
}
|
||||
@@ -748,7 +748,9 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP
|
||||
this.messageStore.removeMessageGroup(correlationKey);
|
||||
}
|
||||
|
||||
protected int findLastReleasedSequenceNumber(Object groupId, Collection<Message<?>> partialSequence) {
|
||||
protected int findLastReleasedSequenceNumber(@SuppressWarnings("unused") Object groupId,
|
||||
Collection<Message<?>> partialSequence) {
|
||||
|
||||
Message<?> lastReleasedMessage = Collections.max(partialSequence, this.sequenceNumberComparator);
|
||||
return new IntegrationMessageHeaderAccessor(lastReleasedMessage).getSequenceNumber();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user