diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaOutboundGatewaySpec.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaOutboundGatewaySpec.java index 38b5da3276..79ae5a7e11 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaOutboundGatewaySpec.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaOutboundGatewaySpec.java @@ -16,6 +16,7 @@ package org.springframework.integration.kafka.dsl; +import java.time.Duration; import java.util.Collections; import java.util.Map; import java.util.function.Consumer; @@ -119,14 +120,32 @@ public class KafkaOutboundGatewaySpec taskScheduler(TaskScheduler scheduler) { + public ReplyingKafkaTemplateSpec taskScheduler(TaskScheduler scheduler) { ((ReplyingKafkaTemplate) this.target).setTaskScheduler(scheduler); return this; } + /** + * Default reply timeout. + * @param replyTimeout the timeout. + * @return the spec. + * @deprecated in favor of {@link #defaultReplyTimeout(Duration)}. + */ + @Deprecated @SuppressWarnings("unchecked") - ReplyingKafkaTemplateSpec replyTimeout(long replyTimeout) { - ((ReplyingKafkaTemplate) this.target).setReplyTimeout(replyTimeout); + public ReplyingKafkaTemplateSpec replyTimeout(long replyTimeout) { + ((ReplyingKafkaTemplate) this.target).setDefaultReplyTimeout(Duration.ofMillis(replyTimeout)); + return this; + } + + /** + * Default reply timeout. + * @param replyTimeout the timeout. + * @return the spec. + */ + @SuppressWarnings("unchecked") + public ReplyingKafkaTemplateSpec defaultReplyTimeout(Duration replyTimeout) { + ((ReplyingKafkaTemplate) this.target).setDefaultReplyTimeout(replyTimeout); return this; } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java index 827ac9fad5..9c57333949 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java @@ -270,6 +270,16 @@ public class KafkaMessageSource extends AbstractMessageSource impl } } + /** + * Get a reference to the configured consumer properties; allows further + * customization of the properties before the source is started. + * @return the properties. + * @since 3.2 + */ + public ConsumerProperties getConsumerProperties() { + return this.consumerProperties; + } + protected String getGroupId() { return this.consumerProperties.getGroupId(); } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java index 0ece863e7a..1a10315eba 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/dsl/KafkaDslTests.java @@ -18,6 +18,7 @@ package org.springframework.integration.kafka.dsl; import static org.assertj.core.api.Assertions.assertThat; +import java.time.Duration; import java.util.Collection; import java.util.Collections; import java.util.Map; @@ -354,7 +355,7 @@ public class KafkaDslTests { return IntegrationFlows.from(Gate.class) .handle(Kafka.outboundGateway(producerFactory(), replyContainer()) .sync(true) - .configureKafkaTemplate(t -> t.replyTimeout(30_000))) + .configureKafkaTemplate(t -> t.defaultReplyTimeout(Duration.ofSeconds(30)))) .get(); } diff --git a/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/KafkaDslKotlinTests.kt b/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/KafkaDslKotlinTests.kt index 0ff59bac09..63aedd9ad4 100644 --- a/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/KafkaDslKotlinTests.kt +++ b/spring-integration-kafka/src/test/kotlin/org/springframework/integration/kafka/dsl/KafkaDslKotlinTests.kt @@ -73,6 +73,7 @@ import org.springframework.messaging.support.GenericMessage import org.springframework.retry.support.RetryTemplate import org.springframework.test.annotation.DirtiesContext import org.springframework.test.context.junit4.SpringRunner +import java.time.Duration import java.util.concurrent.CountDownLatch import java.util.concurrent.TimeUnit import java.util.stream.Stream @@ -328,7 +329,7 @@ class KafkaDslKotlinTests { fun replyingKafkaTemplate() = ReplyingKafkaTemplate(producerFactory(), replyContainer()) .also { - it.setReplyTimeout(30000) + it.setDefaultReplyTimeout(Duration.ofSeconds(30)) } @Bean