From c570bee1979d5dd1fb701c248c6af8cb6b47c3d4 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 25 Apr 2019 14:04:31 -0400 Subject: [PATCH] Fix Sonar complexity issues Address a few complexity issues. --- .../amqp/support/DefaultAmqpHeaderMapper.java | 249 +++++++----------- ...tractAggregatingMessageGroupProcessor.java | 21 +- .../AbstractCorrelatingMessageHandler.java | 8 +- 3 files changed, 108 insertions(+), 170 deletions(-) diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapper.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapper.java index f5bbd1ecc6..1f2364ba92 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapper.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapper.java @@ -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 extractStandardHeaders(MessageProperties amqpMessageProperties) { Map headers = new HashMap(); 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 allHeaders, Map 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 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 amqpMessageProperties.setType(type)) + .acceptIfNotNull(getHeaderIfAvailable(headers, AmqpHeaders.USER_ID, String.class), + userId -> amqpMessageProperties.setUserId(userId)); Map jsonHeaders = new HashMap(); @@ -346,14 +280,11 @@ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper amqpMessageProperties.setHeader("spring_reply_correlation", replyCorrelation)) + .acceptIfHasText(getHeaderIfAvailable(headers, AmqpHeaders.SPRING_REPLY_TO_STACK, String.class), + replyToStack -> amqpMessageProperties.setHeader("spring_reply_to", replyToStack)); } @Override diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractAggregatingMessageGroupProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractAggregatingMessageGroupProcessor.java index 01e946a443..67b7d2e562 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractAggregatingMessageGroupProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractAggregatingMessageGroupProcessor.java @@ -105,6 +105,18 @@ public abstract class AbstractAggregatingMessageGroupProcessor implements Messag */ protected Map aggregateHeaders(MessageGroup group) { Map aggregatedHeaders = new HashMap<>(); + Set 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 doAggregateHeaders(MessageGroup group, Map aggregatedHeaders) { Set conflictKeys = new HashSet<>(); for (Message message : group.getMessages()) { for (Entry 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 defaultHeaders); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java index 7c67be533b..8726cef29c 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java @@ -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> partialSequence) { + protected int findLastReleasedSequenceNumber(@SuppressWarnings("unused") Object groupId, + Collection> partialSequence) { + Message lastReleasedMessage = Collections.max(partialSequence, this.sequenceNumberComparator); return new IntegrationMessageHeaderAccessor(lastReleasedMessage).getSequenceNumber(); }