diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/Kafka.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/Kafka.java index a48e4b1ea7..5dd14bbaa3 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/Kafka.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/Kafka.java @@ -16,6 +16,7 @@ package org.springframework.integration.kafka.dsl; +import java.util.Arrays; import java.util.regex.Pattern; import org.apache.kafka.common.TopicPartition; @@ -29,7 +30,7 @@ import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.listener.ContainerProperties; import org.springframework.kafka.listener.GenericMessageListenerContainer; import org.springframework.kafka.requestreply.ReplyingKafkaTemplate; -import org.springframework.kafka.support.TopicPartitionInitialOffset; +import org.springframework.kafka.support.TopicPartitionOffset; /** * Factory class for Apache Kafka components. @@ -214,6 +215,26 @@ public final class Kafka { containerProperties), listenerMode); } + /** + * Create an initial + * {@link KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec}. + * @param consumerFactory the {@link ConsumerFactory}. + * @param topicPartitions the {@link TopicPartition} vararg. + * @param the Kafka message key type. + * @param the Kafka message value type. + * @return the KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec. + * @deprecated in favor of {@link #messageDrivenChannelAdapter(ConsumerFactory, TopicPartitionOffset...)}. + */ + @Deprecated + public static + KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec messageDrivenChannelAdapter( + ConsumerFactory consumerFactory, + org.springframework.kafka.support.TopicPartitionInitialOffset... topicPartitions) { + + return messageDrivenChannelAdapter(consumerFactory, KafkaMessageDrivenChannelAdapter.ListenerMode.record, + topicPartitions); + } + /** * Create an initial * {@link KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec}. @@ -226,12 +247,39 @@ public final class Kafka { public static KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec messageDrivenChannelAdapter( ConsumerFactory consumerFactory, - TopicPartitionInitialOffset... topicPartitions) { + TopicPartitionOffset... topicPartitions) { return messageDrivenChannelAdapter(consumerFactory, KafkaMessageDrivenChannelAdapter.ListenerMode.record, topicPartitions); } + /** + * Create an initial + * {@link KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec}. + * @param consumerFactory the {@link ConsumerFactory}. + * @param listenerMode the {@link KafkaMessageDrivenChannelAdapter.ListenerMode}. + * @param topicPartitions the {@link TopicPartition} vararg. + * @param the Kafka message key type. + * @param the Kafka message value type. + * @return the + * KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec. + * @deprecated in favor of + * {@link #messageDrivenChannelAdapter(ConsumerFactory, org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter.ListenerMode, TopicPartitionOffset...)} + */ + @Deprecated + public static + KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec messageDrivenChannelAdapter( + ConsumerFactory consumerFactory, + KafkaMessageDrivenChannelAdapter.ListenerMode listenerMode, + org.springframework.kafka.support.TopicPartitionInitialOffset... topicPartitions) { + + return messageDrivenChannelAdapter( + new KafkaMessageListenerContainerSpec<>(consumerFactory, + Arrays.stream(topicPartitions) + .map(org.springframework.kafka.support.TopicPartitionInitialOffset::toTPO) + .toArray(TopicPartitionOffset[]::new)), listenerMode); + } + /** * Create an initial * {@link KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec}. @@ -246,7 +294,7 @@ public final class Kafka { KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec messageDrivenChannelAdapter( ConsumerFactory consumerFactory, KafkaMessageDrivenChannelAdapter.ListenerMode listenerMode, - TopicPartitionInitialOffset... topicPartitions) { + TopicPartitionOffset... topicPartitions) { return messageDrivenChannelAdapter( new KafkaMessageListenerContainerSpec<>(consumerFactory, diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaInboundGatewaySpec.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaInboundGatewaySpec.java index 0f87ec4b67..5976df3b9d 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaInboundGatewaySpec.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaInboundGatewaySpec.java @@ -123,8 +123,7 @@ public class KafkaInboundGatewaySpec the reply value type. */ public static class KafkaInboundGatewayListenerContainerSpec extends - KafkaInboundGatewaySpec> - implements ComponentsRegistration { + KafkaInboundGatewaySpec> { private final KafkaMessageListenerContainerSpec containerSpec; diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageDrivenChannelAdapterSpec.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageDrivenChannelAdapterSpec.java index 949951be76..368e60c55d 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageDrivenChannelAdapterSpec.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageDrivenChannelAdapterSpec.java @@ -149,7 +149,7 @@ public class KafkaMessageDrivenChannelAdapterSpec payloadType) { this.target.setPayloadType(payloadType); return _this(); } @@ -199,8 +199,7 @@ public class KafkaMessageDrivenChannelAdapterSpec the value type. */ public static class KafkaMessageDrivenChannelAdapterListenerContainerSpec extends - KafkaMessageDrivenChannelAdapterSpec> - implements ComponentsRegistration { + KafkaMessageDrivenChannelAdapterSpec> { private final KafkaMessageListenerContainerSpec spec; diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageListenerContainerSpec.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageListenerContainerSpec.java index fa5947a0cb..7b4e869678 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageListenerContainerSpec.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/KafkaMessageListenerContainerSpec.java @@ -30,7 +30,7 @@ import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.kafka.listener.ContainerProperties; import org.springframework.kafka.listener.ErrorHandler; import org.springframework.kafka.listener.GenericErrorHandler; -import org.springframework.kafka.support.TopicPartitionInitialOffset; +import org.springframework.kafka.support.TopicPartitionOffset; /** * A helper class in the Builder pattern style to delegate options to the @@ -54,7 +54,7 @@ public class KafkaMessageListenerContainerSpec } KafkaMessageListenerContainerSpec(ConsumerFactory consumerFactory, - TopicPartitionInitialOffset... topicPartitions) { + TopicPartitionOffset... topicPartitions) { this(consumerFactory, new ContainerProperties(topicPartitions)); } diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java index 06c75f8805..949c10bd16 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java @@ -170,13 +170,13 @@ public class KafkaInboundGateway extends MessagingGatewaySupport implem @Override protected void onInit() { super.onInit(); - MessageListener listener = this.listener; + MessageListener kafkaListener = this.listener; if (this.retryTemplate != null) { - listener = new RetryingMessageListenerAdapter<>(listener, this.retryTemplate, + kafkaListener = new RetryingMessageListenerAdapter<>(kafkaListener, this.retryTemplate, this.recoveryCallback); this.retryTemplate.registerListener(this.listener); } - this.messageListenerContainer.getContainerProperties().setMessageListener(listener); + this.messageListenerContainer.getContainerProperties().setMessageListener(kafkaListener); } @Override diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/InboundGatewayTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/InboundGatewayTests.java index c298017c52..36d2e180c4 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/InboundGatewayTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/InboundGatewayTests.java @@ -96,22 +96,22 @@ public class InboundGatewayTests { @Test public void testInbound() throws Exception { - EmbeddedKafkaBroker embeddedKafka = InboundGatewayTests.embeddedKafka.getEmbeddedKafka(); + EmbeddedKafkaBroker embedded = InboundGatewayTests.embeddedKafka.getEmbeddedKafka(); Map consumerProps = - KafkaTestUtils.consumerProps("replyHandler1", "false", embeddedKafka); + KafkaTestUtils.consumerProps("replyHandler1", "false", embedded); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); ConsumerFactory cf2 = new DefaultKafkaConsumerFactory<>(consumerProps); Consumer consumer = cf2.createConsumer(); - embeddedKafka.consumeFromAnEmbeddedTopic(consumer, topic2); + embedded.consumeFromAnEmbeddedTopic(consumer, topic2); - Map props = KafkaTestUtils.consumerProps("test1", "false", embeddedKafka); + Map props = KafkaTestUtils.consumerProps("test1", "false", embedded); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(props); ContainerProperties containerProps = new ContainerProperties(topic1); containerProps.setIdleEventInterval(100L); KafkaMessageListenerContainer container = new KafkaMessageListenerContainer<>(cf, containerProps); - Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); + Map senderProps = KafkaTestUtils.producerProps(embedded); ProducerFactory pf = new DefaultKafkaProducerFactory<>(senderProps); KafkaTemplate template = new KafkaTemplate<>(pf); template.setDefaultTopic(topic1); @@ -127,8 +127,9 @@ public class InboundGatewayTests { @Override public Message toMessage(ConsumerRecord record, Acknowledgment acknowledgment, - Consumer consumer, Type type) { - Message message = super.toMessage(record, acknowledgment, consumer, type); + Consumer con, Type type) { + + Message message = super.toMessage(record, acknowledgment, con, type); return MessageBuilder.fromMessage(message) .setHeader("testHeader", "testValue") .setHeader(KafkaHeaders.REPLY_TOPIC, topic2) @@ -170,7 +171,7 @@ public class InboundGatewayTests { } @Test - public void testInboundErrorRecover() throws Exception { + public void testInboundErrorRecover() { EmbeddedKafkaBroker broker = InboundGatewayTests.embeddedKafka.getEmbeddedKafka(); Map consumerProps = KafkaTestUtils.consumerProps("replyHandler2", "false", broker); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); @@ -205,8 +206,8 @@ public class InboundGatewayTests { @Override public Message toMessage(ConsumerRecord record, Acknowledgment acknowledgment, - Consumer consumer, Type type) { - Message message = super.toMessage(record, acknowledgment, consumer, type); + Consumer con, Type type) { + Message message = super.toMessage(record, acknowledgment, con, type); return MessageBuilder.fromMessage(message) .setHeader("testHeader", "testValue") .setHeader(KafkaHeaders.REPLY_TOPIC, topic4) @@ -250,21 +251,21 @@ public class InboundGatewayTests { } @Test - public void testInboundRetryErrorRecover() throws Exception { - EmbeddedKafkaBroker embeddedKafka = InboundGatewayTests.embeddedKafka.getEmbeddedKafka(); - Map consumerProps = KafkaTestUtils.consumerProps("replyHandler3", "false", embeddedKafka); + public void testInboundRetryErrorRecover() { + EmbeddedKafkaBroker embedded = InboundGatewayTests.embeddedKafka.getEmbeddedKafka(); + Map consumerProps = KafkaTestUtils.consumerProps("replyHandler3", "false", embedded); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); ConsumerFactory cf2 = new DefaultKafkaConsumerFactory<>(consumerProps); Consumer consumer = cf2.createConsumer(); - embeddedKafka.consumeFromAnEmbeddedTopic(consumer, topic6); + embedded.consumeFromAnEmbeddedTopic(consumer, topic6); - Map props = KafkaTestUtils.consumerProps("test3", "false", embeddedKafka); + Map props = KafkaTestUtils.consumerProps("test3", "false", embedded); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(props); ContainerProperties containerProps = new ContainerProperties(topic5); KafkaMessageListenerContainer container = new KafkaMessageListenerContainer<>(cf, containerProps); - Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); + Map senderProps = KafkaTestUtils.producerProps(embedded); ProducerFactory pf = new DefaultKafkaProducerFactory<>(senderProps); KafkaTemplate template = new KafkaTemplate<>(pf); template.setDefaultTopic(topic5); @@ -285,8 +286,8 @@ public class InboundGatewayTests { @Override public Message toMessage(ConsumerRecord record, Acknowledgment acknowledgment, - Consumer consumer, Type type) { - Message message = super.toMessage(record, acknowledgment, consumer, type); + Consumer con, Type type) { + Message message = super.toMessage(record, acknowledgment, con, type); return MessageBuilder.fromMessage(message) .setHeader("testHeader", "testValue") .setHeader(KafkaHeaders.REPLY_TOPIC, topic6) diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java index 9a35e3a327..9d3af6c38e 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java @@ -69,7 +69,7 @@ import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.DefaultKafkaHeaderMapper; import org.springframework.kafka.support.KafkaHeaders; import org.springframework.kafka.support.KafkaNull; -import org.springframework.kafka.support.TopicPartitionInitialOffset; +import org.springframework.kafka.support.TopicPartitionOffset; import org.springframework.kafka.support.converter.BatchMessageConverter; import org.springframework.kafka.support.converter.BatchMessagingMessageConverter; import org.springframework.kafka.support.converter.ConversionException; @@ -119,7 +119,7 @@ public class MessageDrivenAdapterTests { private static EmbeddedKafkaBroker embeddedKafka = embeddedKafkaRule.getEmbeddedKafka(); @Test - public void testInboundRecord() throws Exception { + public void testInboundRecord() { Map props = KafkaTestUtils.consumerProps("test1", "true", embeddedKafka); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(props); @@ -204,7 +204,7 @@ public class MessageDrivenAdapterTests { } @Test - public void testInboundRecordRetryRecover() throws Exception { + public void testInboundRecordRetryRecover() { Map props = KafkaTestUtils.consumerProps("test4", "true", embeddedKafka); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(props); @@ -257,7 +257,7 @@ public class MessageDrivenAdapterTests { } @Test - public void testInboundRecordNoRetryRecover() throws Exception { + public void testInboundRecordNoRetryRecover() { Map props = KafkaTestUtils.consumerProps("test5", "true", embeddedKafka); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(props); @@ -515,8 +515,8 @@ public class MessageDrivenAdapterTests { resumeLatch.countDown(); return null; }).given(consumer).resume(records.keySet()); - TopicPartitionInitialOffset[] topicPartition = new TopicPartitionInitialOffset[] { - new TopicPartitionInitialOffset("foo", 0) }; + TopicPartitionOffset[] topicPartition = new TopicPartitionOffset[] { + new TopicPartitionOffset("foo", 0) }; ContainerProperties containerProps = new ContainerProperties(topicPartition); containerProps.setAckMode(ContainerProperties.AckMode.RECORD); containerProps.setClientId("clientId");