From e97723484482aa589a7549baad3adfea9c025009 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Fri, 26 Apr 2019 09:22:28 -0400 Subject: [PATCH] Fix new Sonar smells --- .../AbstractAggregatingMessageGroupProcessor.java | 13 ++++++------- .../rsocket/outbound/RSocketOutboundGateway.java | 14 ++++++++------ 2 files changed, 14 insertions(+), 13 deletions(-) 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 67b7d2e562..d0e88bddc0 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 @@ -20,6 +20,7 @@ import java.util.HashMap; import java.util.HashSet; import java.util.Map; import java.util.Map.Entry; +import java.util.Objects; import java.util.Set; import org.apache.commons.logging.Log; @@ -53,11 +54,11 @@ import org.springframework.util.Assert; public abstract class AbstractAggregatingMessageGroupProcessor implements MessageGroupProcessor, BeanFactoryAware { - private final Log logger = LogFactory.getLog(this.getClass()); + protected final Log logger = LogFactory.getLog(getClass()); // NOSONAR - final - private volatile MessageBuilderFactory messageBuilderFactory = new DefaultMessageBuilderFactory(); + private MessageBuilderFactory messageBuilderFactory = new DefaultMessageBuilderFactory(); - private volatile boolean messageBuilderFactorySet; + private boolean messageBuilderFactorySet; private BeanFactory beanFactory; @@ -79,8 +80,7 @@ public abstract class AbstractAggregatingMessageGroupProcessor implements Messag @Override public final Object processMessageGroup(MessageGroup group) { Assert.notNull(group, "MessageGroup must not be null"); - - Map headers = this.aggregateHeaders(group); + Map headers = aggregateHeaders(group); Object payload = this.aggregatePayloads(group, headers); AbstractIntegrationMessageBuilder builder; if (payload instanceof Message) { @@ -131,8 +131,7 @@ public abstract class AbstractAggregatingMessageGroupProcessor implements Messag aggregatedHeaders.put(key, value); } else { - Object existingValue = aggregatedHeaders.get(key); - if (value != existingValue && (value == null || !value.equals(existingValue))) { + if (!Objects.equals(value, aggregatedHeaders.get(key))) { conflictKeys.add(key); } } diff --git a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGateway.java b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGateway.java index 95ffbf0156..c79a7890c5 100644 --- a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGateway.java +++ b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGateway.java @@ -34,9 +34,9 @@ import org.springframework.util.ClassUtils; import org.springframework.util.MimeType; import org.springframework.util.MimeTypeUtils; -import io.rsocket.RSocket; import io.rsocket.RSocketFactory; import io.rsocket.transport.ClientTransport; +import reactor.core.Disposable; import reactor.core.publisher.Mono; /** @@ -167,7 +167,9 @@ public class RSocketOutboundGateway extends AbstractReplyProducingMessageHandler @Override public void destroy() { super.destroy(); - this.rSocketRequesterMono.block().rsocket().dispose(); + this.rSocketRequesterMono.map(RSocketRequester::rsocket) + .doOnNext(Disposable::dispose) + .subscribe(); } @Override @@ -275,20 +277,20 @@ public class RSocketOutboundGateway extends AbstractReplyProducingMessageHandler public enum Command { /** - * Perform {@link RSocket#fireAndForget fireAndForget}. + * Perform {@link io.rsocket.RSocket#fireAndForget fireAndForget}. * @see RSocketRequester.ResponseSpec#send() */ fireAndForget, /** - * Perform {@link RSocket#requestResponse requestResponse}. + * Perform {@link io.rsocket.RSocket#requestResponse requestResponse}. * @see RSocketRequester.ResponseSpec#retrieveMono */ requestResponse, /** - * Perform {@link RSocket#requestStream requestStream} or - * {@link RSocket#requestChannel requestChannel} depending on whether + * Perform {@link io.rsocket.RSocket#requestStream requestStream} or + * {@link io.rsocket.RSocket#requestChannel requestChannel} depending on whether * the request input consists of a single or multiple payloads. * @see RSocketRequester.ResponseSpec#retrieveFlux */