diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java index 84efd7dbd6..b7c342a73f 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java @@ -21,10 +21,16 @@ import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.integration.context.OrderlyShutdownCapable; import org.springframework.integration.endpoint.MessageProducerSupport; import org.springframework.kafka.listener.AbstractMessageListenerContainer; +import org.springframework.kafka.listener.AcknowledgingMessageListener; +import org.springframework.kafka.listener.adapter.FilteringAcknowledgingMessageListenerAdapter; import org.springframework.kafka.listener.adapter.MessagingMessageListenerAdapter; +import org.springframework.kafka.listener.adapter.RecordFilterStrategy; +import org.springframework.kafka.listener.adapter.RetryingAcknowledgingMessageListenerAdapter; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.converter.MessageConverter; import org.springframework.messaging.Message; +import org.springframework.retry.RecoveryCallback; +import org.springframework.retry.support.RetryTemplate; import org.springframework.util.Assert; /** @@ -44,19 +50,115 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo private final MessagingMessageListenerAdapter listener = new IntegrationMessageListener(); + private RecordFilterStrategy recordFilterStrategy; + + private boolean ackDiscarded; + + private RetryTemplate retryTemplate; + + private RecoveryCallback recoveryCallback; + + private boolean filterInRetry; + public KafkaMessageDrivenChannelAdapter(AbstractMessageListenerContainer messageListenerContainer) { Assert.notNull(messageListenerContainer, "messageListenerContainer is required"); Assert.isNull(messageListenerContainer.getContainerProperties().getMessageListener(), "Container must not already have a listener"); this.messageListenerContainer = messageListenerContainer; this.messageListenerContainer.setAutoStartup(false); - this.messageListenerContainer.getContainerProperties().setMessageListener(this.listener); } public void setMessageConverter(MessageConverter messageConverter) { this.listener.setMessageConverter(messageConverter); } + /** + * Specify a {@link RecordFilterStrategy} to wrap + * {@link KafkaMessageDrivenChannelAdapter.IntegrationMessageListener} into + * {@link FilteringAcknowledgingMessageListenerAdapter}. + * @param recordFilterStrategy the {@link RecordFilterStrategy} to use. + * @since 2.0.1 + */ + public void setRecordFilterStrategy(RecordFilterStrategy recordFilterStrategy) { + this.recordFilterStrategy = recordFilterStrategy; + } + + /** + * A {@code boolean} flag to indicate if {@link FilteringAcknowledgingMessageListenerAdapter} + * should acknowledge discarded records or not. + * Does not make sense if {@link #setRecordFilterStrategy(RecordFilterStrategy)} isn't specified. + * @param ackDiscarded true to ack (commit offset for) discarded messages. + * @since 2.0.1 + */ + public void setAckDiscarded(boolean ackDiscarded) { + this.ackDiscarded = ackDiscarded; + } + + /** + * Specify a {@link RetryTemplate} instance to wrap + * {@link KafkaMessageDrivenChannelAdapter.IntegrationMessageListener} into + * {@link RetryingAcknowledgingMessageListenerAdapter}. + * @param retryTemplate the {@link RetryTemplate} to use. + * @since 2.0.1 + */ + public void setRetryTemplate(RetryTemplate retryTemplate) { + this.retryTemplate = retryTemplate; + } + + /** + * A {@link RecoveryCallback} instance for retry operation; + * if null, the exception will be thrown to the container after retries are exhausted. + * Does not make sense if {@link #setRetryTemplate(RetryTemplate)} isn't specified. + * @param recoveryCallback the recovery callback. + * @since 2.0.1 + */ + public void setRecoveryCallback(RecoveryCallback recoveryCallback) { + this.recoveryCallback = recoveryCallback; + } + + /** + * The {@code boolean} flag to specify the order how + * {@link RetryingAcknowledgingMessageListenerAdapter} and + * {@link FilteringAcknowledgingMessageListenerAdapter} are wrapped to each other, + * if both of them are present. + * Does not make sense if only one of {@link RetryTemplate} or + * {@link RecordFilterStrategy} is present, or any. + * @param filterInRetry the order for {@link RetryingAcknowledgingMessageListenerAdapter} and + * {@link FilteringAcknowledgingMessageListenerAdapter} wrapping. Defaults to {@code false}. + * @since 2.0.1 + */ + public void setFilterInRetry(boolean filterInRetry) { + this.filterInRetry = filterInRetry; + } + + @Override + protected void onInit() { + super.onInit(); + + AcknowledgingMessageListener listener = this.listener; + + boolean filterInRetry = this.filterInRetry && this.retryTemplate != null && this.recordFilterStrategy != null; + + if (filterInRetry) { + listener = new FilteringAcknowledgingMessageListenerAdapter<>(listener, this.recordFilterStrategy, + this.ackDiscarded); + listener = new RetryingAcknowledgingMessageListenerAdapter<>(listener, this.retryTemplate, + this.recoveryCallback); + } + else { + if (this.retryTemplate != null) { + listener = new RetryingAcknowledgingMessageListenerAdapter<>(listener, this.retryTemplate, + this.recoveryCallback); + } + if (this.recordFilterStrategy != null) { + listener = new FilteringAcknowledgingMessageListenerAdapter<>(listener, this.recordFilterStrategy, + this.ackDiscarded); + } + } + + this.messageListenerContainer.getContainerProperties().setMessageListener(listener); + } + @Override protected void doStart() { this.messageListenerContainer.start(); diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParserTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParserTests.java index 4f48b6abae..781d07f006 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParserTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/config/xml/KafkaMessageDrivenChannelAdapterParserTests.java @@ -17,6 +17,9 @@ package org.springframework.integration.kafka.config.xml; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; + +import java.util.Collections; import org.junit.Test; import org.junit.runner.RunWith; @@ -24,9 +27,16 @@ import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.integration.channel.NullChannel; import org.springframework.integration.channel.PublishSubscribeChannel; +import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter; import org.springframework.integration.test.util.TestUtils; +import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.listener.KafkaMessageListenerContainer; +import org.springframework.kafka.listener.adapter.FilteringAcknowledgingMessageListenerAdapter; +import org.springframework.kafka.listener.adapter.RecordFilterStrategy; +import org.springframework.kafka.listener.adapter.RetryingAcknowledgingMessageListenerAdapter; +import org.springframework.kafka.listener.config.ContainerProperties; +import org.springframework.retry.support.RetryTemplate; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @@ -60,4 +70,59 @@ public class KafkaMessageDrivenChannelAdapterParserTests { assertThat(container).isNotNull(); } + @Test + @SuppressWarnings("unchecked") + public void testKafkaMessageDrivenChannelAdapterOptions() { + DefaultKafkaConsumerFactory cf = + new DefaultKafkaConsumerFactory<>(Collections.emptyMap()); + ContainerProperties containerProps = new ContainerProperties("foo"); + KafkaMessageListenerContainer container = + new KafkaMessageListenerContainer<>(cf, containerProps); + KafkaMessageDrivenChannelAdapter adapter = new KafkaMessageDrivenChannelAdapter<>(container); + adapter.setOutputChannel(new QueueChannel()); + + adapter.setRecordFilterStrategy(mock(RecordFilterStrategy.class)); + adapter.afterPropertiesSet(); + + containerProps = TestUtils.getPropertyValue(container, "containerProperties", ContainerProperties.class); + + Object messageListener = containerProps.getMessageListener(); + assertThat(messageListener).isInstanceOf(FilteringAcknowledgingMessageListenerAdapter.class); + + Object delegate = TestUtils.getPropertyValue(messageListener, "delegate"); + + assertThat(delegate.getClass().getName()).contains("$IntegrationMessageListener"); + + adapter.setRecordFilterStrategy(null); + adapter.setRetryTemplate(new RetryTemplate()); + adapter.afterPropertiesSet(); + + messageListener = containerProps.getMessageListener(); + assertThat(messageListener).isInstanceOf(RetryingAcknowledgingMessageListenerAdapter.class); + + delegate = TestUtils.getPropertyValue(messageListener, "delegate"); + + assertThat(delegate.getClass().getName()).contains("$IntegrationMessageListener"); + + adapter.setRecordFilterStrategy(mock(RecordFilterStrategy.class)); + adapter.afterPropertiesSet(); + + messageListener = containerProps.getMessageListener(); + assertThat(messageListener).isInstanceOf(FilteringAcknowledgingMessageListenerAdapter.class); + + delegate = TestUtils.getPropertyValue(messageListener, "delegate"); + + assertThat(delegate).isInstanceOf(RetryingAcknowledgingMessageListenerAdapter.class); + + adapter.setFilterInRetry(true); + adapter.afterPropertiesSet(); + + messageListener = containerProps.getMessageListener(); + assertThat(messageListener).isInstanceOf(RetryingAcknowledgingMessageListenerAdapter.class); + + delegate = TestUtils.getPropertyValue(messageListener, "delegate"); + + assertThat(delegate).isInstanceOf(FilteringAcknowledgingMessageListenerAdapter.class); + } + }