> queue;
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/IntegrationManagementConfigurer.java b/spring-integration-core/src/main/java/org/springframework/integration/config/IntegrationManagementConfigurer.java
index 7b8bf62203..b11293c84a 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/config/IntegrationManagementConfigurer.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/config/IntegrationManagementConfigurer.java
@@ -19,8 +19,6 @@ package org.springframework.integration.config;
import java.util.Map;
import java.util.Map.Entry;
-import org.apache.commons.logging.Log;
-
import org.springframework.beans.BeansException;
import org.springframework.beans.factory.BeanNameAware;
import org.springframework.beans.factory.SmartInitializingSingleton;
@@ -54,7 +52,7 @@ public class IntegrationManagementConfigurer
implements SmartInitializingSingleton, ApplicationContextAware, BeanNameAware, BeanPostProcessor {
/**
- * Bean name of tehe configurer.
+ * Bean name of the configurer.
*/
public static final String MANAGEMENT_CONFIGURER_NAME = "integrationManagementConfigurer";
@@ -86,8 +84,8 @@ public class IntegrationManagementConfigurer
* Exception logging (debug or otherwise) is not affected by this setting.
*
* It has been found that in high-volume messaging environments, calls to methods such as
- * {@link Log#isDebugEnabled()} can be quite expensive and account for an inordinate amount of CPU
- * time.
+ * {@link org.apache.commons.logging.Log#isDebugEnabled()} can be quite expensive
+ * and account for an inordinate amount of CPU time.
*
* Set this to false to disable logging by default in all framework components that implement
* {@link IntegrationManagement} (channels, message handlers etc). This turns off logging such as
@@ -121,7 +119,6 @@ public class IntegrationManagementConfigurer
if (!getOverrides(bean).loggingConfigured) {
bean.setLoggingEnabled(this.defaultLoggingEnabled);
}
- String name = entry.getKey();
}
this.singletonsInstantiated = true;
}
@@ -140,10 +137,8 @@ public class IntegrationManagementConfigurer
@Override
public Object postProcessAfterInitialization(Object bean, String name) throws BeansException {
- if (this.singletonsInstantiated) {
- if (this.metricsCaptor != null && bean instanceof IntegrationManagement) {
- ((IntegrationManagement) bean).registerMetricsCaptor(this.metricsCaptor);
- }
+ if (this.singletonsInstantiated && this.metricsCaptor != null && bean instanceof IntegrationManagement) {
+ ((IntegrationManagement) bean).registerMetricsCaptor(this.metricsCaptor);
}
return bean;
}
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractEndpoint.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractEndpoint.java
index f4b940e416..a8595ca50a 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractEndpoint.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractEndpoint.java
@@ -21,7 +21,6 @@ import java.util.concurrent.locks.ReentrantLock;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.NoSuchBeanDefinitionException;
-import org.springframework.context.SmartLifecycle;
import org.springframework.integration.context.IntegrationContextUtils;
import org.springframework.integration.context.IntegrationObjectSupport;
import org.springframework.integration.context.IntegrationProperties;
@@ -80,7 +79,7 @@ public abstract class AbstractEndpoint extends IntegrationObjectSupport
* Such endpoints can be started/stopped as a group.
* @param role the role for this endpoint.
* @since 5.0
- * @see SmartLifecycle
+ * @see org.springframework.context.SmartLifecycle
* @see org.springframework.integration.support.SmartLifecycleRoleController
*/
public void setRole(String role) {
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/graph/MessageChannelNode.java b/spring-integration-core/src/main/java/org/springframework/integration/graph/MessageChannelNode.java
index 0d07cfc76e..0e7412e56d 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/graph/MessageChannelNode.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/graph/MessageChannelNode.java
@@ -30,7 +30,6 @@ import org.springframework.messaging.MessageChannel;
* @since 4.3
*
*/
-@SuppressWarnings("deprecation")
public class MessageChannelNode extends IntegrationNode implements SendTimersAware {
private Supplier sendTimers;
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/leader/LockRegistryLeaderInitiator.java b/spring-integration-core/src/main/java/org/springframework/integration/support/leader/LockRegistryLeaderInitiator.java
index f3f5718434..ef98dd73b1 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/support/leader/LockRegistryLeaderInitiator.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/support/leader/LockRegistryLeaderInitiator.java
@@ -36,7 +36,6 @@ import org.springframework.integration.leader.DefaultCandidate;
import org.springframework.integration.leader.event.DefaultLeaderEventPublisher;
import org.springframework.integration.leader.event.LeaderEventPublisher;
import org.springframework.integration.support.locks.LockRegistry;
-import org.springframework.integration.support.management.ManageableSmartLifecycle;
import org.springframework.scheduling.concurrent.CustomizableThreadFactory;
import org.springframework.util.Assert;
@@ -61,7 +60,7 @@ import org.springframework.util.Assert;
*
* @since 4.3.1
*/
-public class LockRegistryLeaderInitiator implements ManageableSmartLifecycle, DisposableBean,
+public class LockRegistryLeaderInitiator implements SmartLifecycle, DisposableBean,
ApplicationEventPublisherAware {
public static final long DEFAULT_HEART_BEAT_TIME = 500L;
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagement.java b/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagement.java
index 765e40690f..600ca1e52a 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagement.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagement.java
@@ -20,7 +20,6 @@ import org.springframework.beans.factory.DisposableBean;
import org.springframework.integration.support.context.NamedComponent;
import org.springframework.integration.support.management.metrics.MetricsCaptor;
import org.springframework.jmx.export.annotation.ManagedAttribute;
-import org.springframework.lang.Nullable;
/**
* Base interface for Integration managed components.
@@ -54,7 +53,6 @@ public interface IntegrationManagement extends NamedComponent, DisposableBean {
default void setManagedName(String managedName) {
}
- @Nullable
default String getManagedName() {
return null;
}
@@ -62,7 +60,6 @@ public interface IntegrationManagement extends NamedComponent, DisposableBean {
default void setManagedType(String managedType) {
}
- @Nullable
default String getManagedType() {
return null;
}
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/management/MessageSourceManagement.java b/spring-integration-core/src/main/java/org/springframework/integration/support/management/MessageSourceManagement.java
index 87f7099218..9f13f3c312 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/support/management/MessageSourceManagement.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/support/management/MessageSourceManagement.java
@@ -26,7 +26,6 @@ import org.springframework.jmx.export.annotation.ManagedAttribute;
* @since 5.0
*
*/
-@SuppressWarnings("deprecation")
@IntegrationManagedResource
public interface MessageSourceManagement {
diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java
index 5f41d134ac..7173b74767 100644
--- a/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java
+++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveStreamsConsumerTests.java
@@ -26,6 +26,7 @@ import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
+import java.time.Duration;
import java.util.LinkedList;
import java.util.List;
import java.util.concurrent.BlockingQueue;
@@ -329,7 +330,7 @@ public class ReactiveStreamsConsumerTests {
StepVerifier.create(sink.asFlux())
.expectNext(testMessage, testMessage2)
.thenCancel()
- .verify();
+ .verify(Duration.ofSeconds(10));
reactiveConsumer.stop();
}
diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ReactiveInboundChannelAdapterTests.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ReactiveInboundChannelAdapterTests.java
index 6e50f8c759..e2c42df751 100644
--- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ReactiveInboundChannelAdapterTests.java
+++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ReactiveInboundChannelAdapterTests.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2018-2019 the original author or authors.
+ * Copyright 2018-2020 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
@@ -16,6 +16,7 @@
package org.springframework.integration.endpoint;
+import java.time.Duration;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Supplier;
@@ -60,7 +61,7 @@ public class ReactiveInboundChannelAdapterTests {
StepVerifier.create(testFlux)
.expectNext(2, 4, 6, 8, 10, 12, 14, 16)
.thenCancel()
- .verify();
+ .verify(Duration.ofSeconds(10));
}
@Configuration
diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ReactiveMessageProducerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ReactiveMessageProducerTests.java
index fff06a4c5d..ffe8f15ed9 100644
--- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ReactiveMessageProducerTests.java
+++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ReactiveMessageProducerTests.java
@@ -18,6 +18,8 @@ package org.springframework.integration.endpoint;
import static org.assertj.core.api.Assertions.assertThat;
+import java.time.Duration;
+
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
@@ -58,7 +60,7 @@ public class ReactiveMessageProducerTests {
.cast(String.class))
.expectNext("test1", "test2")
.thenCancel()
- .verify();
+ .verify(Duration.ofSeconds(10));
assertThat(this.producer.isRunning()).isFalse();
}
diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ReactiveMessageSourceProducerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ReactiveMessageSourceProducerTests.java
index 3fbe2e8f3d..f5831a250c 100644
--- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ReactiveMessageSourceProducerTests.java
+++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/ReactiveMessageSourceProducerTests.java
@@ -97,7 +97,7 @@ public class ReactiveMessageSourceProducerTests {
reactiveMessageSourceProducer.start();
- stepVerifier.verify();
+ stepVerifier.verify(Duration.ofSeconds(10));
reactiveMessageSourceProducer.stop();
diff --git a/spring-integration-core/src/test/kotlin/org/springframework/integration/dsl/KotlinDslTests.kt b/spring-integration-core/src/test/kotlin/org/springframework/integration/dsl/KotlinDslTests.kt
index 3bf29a29e2..f7bc283e6d 100644
--- a/spring-integration-core/src/test/kotlin/org/springframework/integration/dsl/KotlinDslTests.kt
+++ b/spring-integration-core/src/test/kotlin/org/springframework/integration/dsl/KotlinDslTests.kt
@@ -47,6 +47,7 @@ import org.springframework.test.annotation.DirtiesContext
import org.springframework.test.context.junit.jupiter.SpringJUnitConfig
import reactor.core.publisher.Flux
import reactor.test.StepVerifier
+import java.time.Duration
import java.util.*
import java.util.concurrent.atomic.AtomicReference
import java.util.function.Function
@@ -167,7 +168,7 @@ class KotlinDslTests {
val registration = this.integrationFlowContext.registration(integrationFlow).register()
- verifyLater.verify()
+ verifyLater.verify(Duration.ofSeconds(10))
registration.destroy()
}
diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/remote/aop/RotatingServerAdvice.java b/spring-integration-file/src/main/java/org/springframework/integration/file/remote/aop/RotatingServerAdvice.java
index 47088d9ea5..79c276fb0b 100644
--- a/spring-integration-file/src/main/java/org/springframework/integration/file/remote/aop/RotatingServerAdvice.java
+++ b/spring-integration-file/src/main/java/org/springframework/integration/file/remote/aop/RotatingServerAdvice.java
@@ -36,10 +36,7 @@ import org.springframework.util.Assert;
* @since 5.0.7
*
*/
-@SuppressWarnings("deprecation")
-public class RotatingServerAdvice
- extends org.springframework.integration.aop.AbstractMessageSourceAdvice
- implements MessageSourceMutator {
+public class RotatingServerAdvice implements MessageSourceMutator {
private final RotationPolicy rotationPolicy;
diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpOutboundGateway.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpOutboundGateway.java
index c48db1c7e9..cf3525c2da 100644
--- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpOutboundGateway.java
+++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/TcpOutboundGateway.java
@@ -25,7 +25,6 @@ import java.util.concurrent.Semaphore;
import java.util.concurrent.TimeUnit;
import org.springframework.context.ApplicationEventPublisher;
-import org.springframework.context.Lifecycle;
import org.springframework.expression.EvaluationContext;
import org.springframework.expression.Expression;
import org.springframework.expression.common.LiteralExpression;
@@ -58,7 +57,7 @@ import org.springframework.util.concurrent.SettableListenableFuture;
* (or times out). Asynchronous requests/responses over the same connection are not
* supported - use a pair of outbound/inbound adapters for that use case.
*
- * {@link Lifecycle} methods delegate to the underlying {@link AbstractConnectionFactory}
+ * {@link org.springframework.context.Lifecycle} methods delegate to the underlying {@link AbstractConnectionFactory}.
*
*
* @author Gary Russell
diff --git a/spring-integration-jmx/src/main/java/org/springframework/integration/jmx/config/IntegrationMBeanExportConfiguration.java b/spring-integration-jmx/src/main/java/org/springframework/integration/jmx/config/IntegrationMBeanExportConfiguration.java
index 4e0ad334fa..74f331cd7c 100644
--- a/spring-integration-jmx/src/main/java/org/springframework/integration/jmx/config/IntegrationMBeanExportConfiguration.java
+++ b/spring-integration-jmx/src/main/java/org/springframework/integration/jmx/config/IntegrationMBeanExportConfiguration.java
@@ -26,8 +26,6 @@ import javax.management.MBeanServer;
import org.springframework.beans.BeansException;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.beans.factory.BeanFactoryAware;
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.config.BeanDefinition;
import org.springframework.beans.factory.config.BeanExpressionContext;
import org.springframework.beans.factory.config.BeanExpressionResolver;
@@ -42,7 +40,6 @@ import org.springframework.context.expression.StandardBeanExpressionResolver;
import org.springframework.core.annotation.AnnotationAttributes;
import org.springframework.core.env.Environment;
import org.springframework.core.type.AnnotationMetadata;
-import org.springframework.integration.config.IntegrationManagementConfigurer;
import org.springframework.integration.monitor.IntegrationMBeanExporter;
import org.springframework.util.Assert;
import org.springframework.util.StringUtils;
@@ -76,10 +73,6 @@ public class IntegrationMBeanExportConfiguration implements ImportAware, Environ
private Environment environment;
- @Autowired(required = false)
- @Qualifier(IntegrationManagementConfigurer.MANAGEMENT_CONFIGURER_NAME)
- private IntegrationManagementConfigurer configurer;
-
@Override
public void setBeanFactory(BeanFactory beanFactory) throws BeansException {
this.beanFactory = beanFactory;
@@ -102,7 +95,6 @@ public class IntegrationMBeanExportConfiguration implements ImportAware, Environ
"@EnableIntegrationMBeanExport is not present on importing class " + importMetadata.getClassName());
}
- @SuppressWarnings("deprecation")
@Bean(name = MBEAN_EXPORTER_NAME)
@Role(BeanDefinition.ROLE_INFRASTRUCTURE)
public IntegrationMBeanExporter mbeanExporter() {
diff --git a/spring-integration-jmx/src/main/java/org/springframework/integration/monitor/IntegrationMBeanExporter.java b/spring-integration-jmx/src/main/java/org/springframework/integration/monitor/IntegrationMBeanExporter.java
index 25839d24aa..3a109a3f3f 100644
--- a/spring-integration-jmx/src/main/java/org/springframework/integration/monitor/IntegrationMBeanExporter.java
+++ b/spring-integration-jmx/src/main/java/org/springframework/integration/monitor/IntegrationMBeanExporter.java
@@ -107,7 +107,6 @@ import org.springframework.util.ReflectionUtils;
* @author Meherzad Lahewala
*/
@ManagedResource
-@SuppressWarnings("deprecation")
public class IntegrationMBeanExporter extends MBeanExporter
implements ApplicationContextAware, DestructionAwareBeanPostProcessor {
@@ -148,9 +147,7 @@ public class IntegrationMBeanExporter extends MBeanExporter
private String domain = DEFAULT_DOMAIN;
- private String[] componentNamePatterns = { "*" };
-
- private IntegrationManagementConfigurer managementConfigurer;
+ private String[] componentNamePatterns = {"*"};
private volatile long shutdownDeadline;
@@ -247,7 +244,7 @@ public class IntegrationMBeanExporter extends MBeanExporter
private void populateMessageHandlers() {
Map messageHandlers = this.applicationContext
- .getBeansOfType(MessageHandler.class);
+ .getBeansOfType(MessageHandler.class);
for (Entry entry : messageHandlers.entrySet()) {
String beanName = entry.getKey();
@@ -269,7 +266,7 @@ public class IntegrationMBeanExporter extends MBeanExporter
private void populateMessageSources() {
this.applicationContext.getBeansOfType(
- IntegrationInboundManagement.class)
+ IntegrationInboundManagement.class)
.values()
.stream()
// If the source is proxied, we have to extract the target to expose as an MBean.
@@ -301,15 +298,10 @@ public class IntegrationMBeanExporter extends MBeanExporter
private void configureManagementConfigurer() {
if (!this.applicationContext.containsBean(IntegrationManagementConfigurer.MANAGEMENT_CONFIGURER_NAME)) {
- this.managementConfigurer = new IntegrationManagementConfigurer();
- this.managementConfigurer.setApplicationContext(this.applicationContext);
- this.managementConfigurer.setBeanName(IntegrationManagementConfigurer.MANAGEMENT_CONFIGURER_NAME);
- this.managementConfigurer.afterSingletonsInstantiated();
- }
- else {
- this.managementConfigurer =
- this.applicationContext.getBean(IntegrationManagementConfigurer.MANAGEMENT_CONFIGURER_NAME,
- IntegrationManagementConfigurer.class);
+ IntegrationManagementConfigurer managementConfigurer = new IntegrationManagementConfigurer();
+ managementConfigurer.setApplicationContext(this.applicationContext);
+ managementConfigurer.setBeanName(IntegrationManagementConfigurer.MANAGEMENT_CONFIGURER_NAME);
+ managementConfigurer.afterSingletonsInstantiated();
}
}
diff --git a/spring-integration-jmx/src/main/java/org/springframework/integration/monitor/ManagedEndpoint.java b/spring-integration-jmx/src/main/java/org/springframework/integration/monitor/ManagedEndpoint.java
index 2fa008068b..1abb43215a 100644
--- a/spring-integration-jmx/src/main/java/org/springframework/integration/monitor/ManagedEndpoint.java
+++ b/spring-integration-jmx/src/main/java/org/springframework/integration/monitor/ManagedEndpoint.java
@@ -19,7 +19,6 @@ package org.springframework.integration.monitor;
import org.springframework.context.Lifecycle;
import org.springframework.integration.endpoint.AbstractEndpoint;
import org.springframework.integration.support.management.IntegrationManagedResource;
-import org.springframework.integration.support.management.ManageableLifecycle;
import org.springframework.jmx.export.annotation.ManagedAttribute;
import org.springframework.jmx.export.annotation.ManagedOperation;
@@ -30,7 +29,7 @@ import org.springframework.jmx.export.annotation.ManagedOperation;
* @author Gary Russell
*
* @deprecated this is no longer used by the framework. Replaced by
- * {@link ManageableLifecycle}.
+ * {@link org.springframework.integration.support.management.ManageableLifecycle}.
*
*/
@Deprecated
diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/channel/PublishSubscribeKafkaChannel.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/channel/PublishSubscribeKafkaChannel.java
index 8c1ae3989c..2450bf4b33 100644
--- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/channel/PublishSubscribeKafkaChannel.java
+++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/channel/PublishSubscribeKafkaChannel.java
@@ -27,7 +27,7 @@ import org.springframework.kafka.core.KafkaOperations;
*
* @author Gary Russell
*
- * @since 4.4
+ * @since 5.4
*
*/
public class PublishSubscribeKafkaChannel extends SubscribableKafkaChannel implements BroadcastCapableChannel {
@@ -40,6 +40,7 @@ public class PublishSubscribeKafkaChannel extends SubscribableKafkaChannel imple
*/
public PublishSubscribeKafkaChannel(KafkaOperations, ?> template, KafkaListenerContainerFactory> factory,
String channelTopic) {
+
super(template, factory, channelTopic);
}
diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/channel/SubscribableKafkaChannel.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/channel/SubscribableKafkaChannel.java
index 0f5ae2865f..0dabc252ea 100644
--- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/channel/SubscribableKafkaChannel.java
+++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/channel/SubscribableKafkaChannel.java
@@ -19,7 +19,6 @@ package org.springframework.integration.kafka.channel;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecord;
-import org.springframework.context.SmartLifecycle;
import org.springframework.integration.dispatcher.MessageDispatcher;
import org.springframework.integration.dispatcher.RoundRobinLoadBalancingStrategy;
import org.springframework.integration.dispatcher.UnicastingDispatcher;
@@ -95,7 +94,7 @@ public class SubscribableKafkaChannel extends AbstractKafkaChannel implements Su
/**
* Set the auto startup.
* @param autoStartup true to automatically start.
- * @see SmartLifecycle
+ * @see org.springframework.context.SmartLifecycle
*/
public void setAutoStartup(boolean autoStartup) {
this.autoStartup = autoStartup;
@@ -123,7 +122,7 @@ public class SubscribableKafkaChannel extends AbstractKafkaChannel implements Su
.dispatch(toMessagingMessage(record, acknowledgment, consumer));
}
- });
+ });
}
protected MessageDispatcher createDispatcher() {
diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/inbound/ReactiveMongoDbMessageSourceTests.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/inbound/ReactiveMongoDbMessageSourceTests.java
index 6e2b042375..4155aa01ed 100644
--- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/inbound/ReactiveMongoDbMessageSourceTests.java
+++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/inbound/ReactiveMongoDbMessageSourceTests.java
@@ -184,7 +184,7 @@ public class ReactiveMongoDbMessageSourceTests extends MongoDbAvailableTests {
.assertNext(
message -> assertThat(((Person) message.getPayload()).getName()).isEqualTo("Oleg"))
.thenCancel()
- .verify();
+ .verify(Duration.ofSeconds(10));
context.close();
}
diff --git a/spring-integration-r2dbc/src/test/java/org/springframework/integration/r2dbc/inbound/R2dbcMessageSourceTests.java b/spring-integration-r2dbc/src/test/java/org/springframework/integration/r2dbc/inbound/R2dbcMessageSourceTests.java
index d9d01b8063..6ca4b88068 100644
--- a/spring-integration-r2dbc/src/test/java/org/springframework/integration/r2dbc/inbound/R2dbcMessageSourceTests.java
+++ b/spring-integration-r2dbc/src/test/java/org/springframework/integration/r2dbc/inbound/R2dbcMessageSourceTests.java
@@ -19,6 +19,7 @@ package org.springframework.integration.r2dbc.inbound;
import static org.assertj.core.api.Assertions.assertThat;
+import java.time.Duration;
import java.util.Arrays;
import java.util.List;
@@ -134,7 +135,7 @@ public class R2dbcMessageSourceTests {
StepVerifier.create((Flux>) r2dbcMessageSourceError.receive().getPayload())
.expectErrorMatches(throwable -> throwable instanceof IllegalStateException
&& throwable.getMessage().contains("'queryExpression' must evaluate to String or"))
- .verify();
+ .verify(Duration.ofSeconds(10));
}
@Configuration
diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducer.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducer.java
index c1a3b49447..285acf070d 100644
--- a/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducer.java
+++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducer.java
@@ -189,11 +189,11 @@ public class ReactiveRedisStreamMessageProducer extends MessageProducerSupport {
Mono> consumerGroupMono = Mono.empty();
if (this.createConsumerGroup) {
consumerGroupMono =
- this.reactiveStreamOperations.createGroup(this.streamKey, this.consumerGroup)
+ this.reactiveStreamOperations.createGroup(this.streamKey, this.consumerGroup) // NOSONAR
.onErrorReturn(this.consumerGroup);
}
- Consumer consumer = Consumer.from(this.consumerGroup, this.consumerName);
+ Consumer consumer = Consumer.from(this.consumerGroup, this.consumerName); // NOSONAR
if (offset.getOffset().equals(ReadOffset.latest())) {
// for consumer group offset id should be equal '>'
diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducerTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducerTests.java
index 3ddb201606..f6b1cbfe9b 100644
--- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducerTests.java
+++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducerTests.java
@@ -100,7 +100,7 @@ public class ReactiveRedisStreamMessageProducerTests extends RedisAvailableTests
.assertNext((infoGroup) ->
assertThat(infoGroup.groupName()).isEqualTo(this.redisStreamMessageProducer.getBeanName()))
.thenCancel()
- .verify();
+ .verify(Duration.ofSeconds(10));
}
@Test
diff --git a/spring-integration-rsocket/src/test/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGatewayIntegrationTests.java b/spring-integration-rsocket/src/test/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGatewayIntegrationTests.java
index d9ff7f8a80..9d59498a15 100644
--- a/spring-integration-rsocket/src/test/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGatewayIntegrationTests.java
+++ b/spring-integration-rsocket/src/test/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGatewayIntegrationTests.java
@@ -160,7 +160,7 @@ public class RSocketOutboundGatewayIntegrationTests {
StepVerifier.create(controller.fireForgetPayloads.asFlux())
.expectNext("Hello")
.thenCancel()
- .verify();
+ .verify(Duration.ofSeconds(10));
disposable.dispose();
}
@@ -195,7 +195,7 @@ public class RSocketOutboundGatewayIntegrationTests {
.setHeader(RSocketRequesterMethodArgumentResolver.RSOCKET_REQUESTER_HEADER, rsocketRequester)
.build());
- verifier.verify();
+ verifier.verify(Duration.ofSeconds(10));
}
@Test
@@ -227,7 +227,7 @@ public class RSocketOutboundGatewayIntegrationTests {
.setHeader(RSocketRequesterMethodArgumentResolver.RSOCKET_REQUESTER_HEADER, rsocketRequester)
.build());
- verifier.verify();
+ verifier.verify(Duration.ofSeconds(10));
}
@Test
@@ -261,7 +261,7 @@ public class RSocketOutboundGatewayIntegrationTests {
.setHeader(RSocketRequesterMethodArgumentResolver.RSOCKET_REQUESTER_HEADER, rsocketRequester)
.build());
- verifier.verify();
+ verifier.verify(Duration.ofSeconds(10));
}
@Test
@@ -295,7 +295,7 @@ public class RSocketOutboundGatewayIntegrationTests {
.setHeader(RSocketRequesterMethodArgumentResolver.RSOCKET_REQUESTER_HEADER, rsocketRequester)
.build());
- verifier.verify();
+ verifier.verify(Duration.ofSeconds(10));
}
@@ -326,7 +326,7 @@ public class RSocketOutboundGatewayIntegrationTests {
.setHeader(RSocketRequesterMethodArgumentResolver.RSOCKET_REQUESTER_HEADER, rsocketRequester)
.build());
- verifier.verify();
+ verifier.verify(Duration.ofSeconds(10));
}
@Test
@@ -356,7 +356,7 @@ public class RSocketOutboundGatewayIntegrationTests {
.setHeader(RSocketRequesterMethodArgumentResolver.RSOCKET_REQUESTER_HEADER, rsocketRequester)
.build());
- verifier.verify();
+ verifier.verify(Duration.ofSeconds(10));
}
@Test
@@ -388,7 +388,7 @@ public class RSocketOutboundGatewayIntegrationTests {
.setHeader(RSocketRequesterMethodArgumentResolver.RSOCKET_REQUESTER_HEADER, rsocketRequester)
.build());
- verifier.verify();
+ verifier.verify(Duration.ofSeconds(10));
}
@Test
@@ -420,7 +420,7 @@ public class RSocketOutboundGatewayIntegrationTests {
.setHeader(RSocketRequesterMethodArgumentResolver.RSOCKET_REQUESTER_HEADER, rsocketRequester)
.build());
- verifier.verify();
+ verifier.verify(Duration.ofSeconds(10));
}
@Test