From ee6b3f87b3bc3df460522676efa0bbe58bf62869 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 25 Apr 2017 11:45:48 -0400 Subject: [PATCH] GH-164: Spring Kafka 2.0.0 Compatibility Resolves: https://github.com/spring-projects/spring-integration-kafka/issues/164 Also gradle 3.5. Requires https://github.com/spring-projects/spring-kafka/pull/296 Fix javadocs More javadoc polishing Updates for new Consumer header * Simple polishing --- .../KafkaMessageDrivenChannelAdapterSpec.java | 31 +++++----- .../KafkaMessageDrivenChannelAdapter.java | 58 ++++++++++--------- ...essageDrivenChannelAdapterParserTests.java | 16 ++--- .../inbound/MessageDrivenAdapterTests.java | 16 +++-- 4 files changed, 64 insertions(+), 57 deletions(-) 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 a7bf6b7ab1..177d6f92f1 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 @@ -33,9 +33,7 @@ import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.listener.AcknowledgingMessageListener; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.kafka.listener.ErrorHandler; -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.kafka.support.TopicPartitionInitialOffset; import org.springframework.kafka.support.converter.BatchMessageConverter; @@ -99,7 +97,7 @@ public class KafkaMessageDrivenChannelAdapterSpec extends MessageProducerSuppo private RetryTemplate retryTemplate; - private RecoveryCallback recoveryCallback; + private RecoveryCallback recoveryCallback; private boolean filterInRetry; @@ -137,7 +138,7 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo /** * Specify a {@link RecordFilterStrategy} to wrap * {@link KafkaMessageDrivenChannelAdapter.IntegrationRecordMessageListener} into - * {@link FilteringAcknowledgingMessageListenerAdapter}. + * {@link FilteringMessageListenerAdapter}. * @param recordFilterStrategy the {@link RecordFilterStrategy} to use. * @since 2.0.1 */ @@ -146,7 +147,7 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo } /** - * A {@code boolean} flag to indicate if {@link FilteringAcknowledgingMessageListenerAdapter} + * A {@code boolean} flag to indicate if {@link FilteringMessageListenerAdapter} * 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. @@ -159,7 +160,7 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo /** * Specify a {@link RetryTemplate} instance to wrap * {@link KafkaMessageDrivenChannelAdapter.IntegrationRecordMessageListener} into - * {@link RetryingAcknowledgingMessageListenerAdapter}. + * {@link RetryingMessageListenerAdapter}. * @param retryTemplate the {@link RetryTemplate} to use. * @since 2.0.1 */ @@ -176,19 +177,19 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo * @param recoveryCallback the recovery callback. * @since 2.0.1 */ - public void setRecoveryCallback(RecoveryCallback recoveryCallback) { + 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, + * {@link RetryingMessageListenerAdapter} and + * {@link FilteringMessageListenerAdapter} 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}. + * @param filterInRetry the order for {@link RetryingMessageListenerAdapter} and + * {@link FilteringMessageListenerAdapter} wrapping. Defaults to {@code false}. * @since 2.0.1 */ public void setFilterInRetry(boolean filterInRetry) { @@ -211,34 +212,34 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo super.onInit(); if (this.mode.equals(ListenerMode.record)) { - AcknowledgingMessageListener listener = this.recordListener; + MessageListener listener = this.recordListener; boolean filterInRetry = this.filterInRetry && this.retryTemplate != null && this.recordFilterStrategy != null; if (filterInRetry) { - listener = new FilteringAcknowledgingMessageListenerAdapter<>(listener, this.recordFilterStrategy, + listener = new FilteringMessageListenerAdapter<>(listener, this.recordFilterStrategy, this.ackDiscarded); - listener = new RetryingAcknowledgingMessageListenerAdapter<>(listener, this.retryTemplate, - this.recoveryCallback); + listener = new RetryingMessageListenerAdapter<>(listener, this.retryTemplate, + this.recoveryCallback); } else { if (this.retryTemplate != null) { - listener = new RetryingAcknowledgingMessageListenerAdapter<>(listener, this.retryTemplate, + listener = new RetryingMessageListenerAdapter<>(listener, this.retryTemplate, this.recoveryCallback); } if (this.recordFilterStrategy != null) { - listener = new FilteringAcknowledgingMessageListenerAdapter<>(listener, this.recordFilterStrategy, + listener = new FilteringMessageListenerAdapter<>(listener, this.recordFilterStrategy, this.ackDiscarded); } } this.messageListenerContainer.getContainerProperties().setMessageListener(listener); } else { - BatchAcknowledgingMessageListener listener = this.batchListener; + BatchMessageListener listener = this.batchListener; if (this.recordFilterStrategy != null) { - listener = new FilteringBatchAcknowledgingMessageListenerAdapter<>(listener, this.recordFilterStrategy, + listener = new FilteringBatchMessageListenerAdapter<>(listener, this.recordFilterStrategy, this.ackDiscarded); } this.messageListenerContainer.getContainerProperties().setMessageListener(listener); @@ -297,10 +298,10 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo } @Override - public void onMessage(ConsumerRecord record, Acknowledgment acknowledgment) { + public void onMessage(ConsumerRecord record, Acknowledgment acknowledgment, Consumer consumer) { Message message = null; try { - message = toMessagingMessage(record, acknowledgment); + message = toMessagingMessage(record, acknowledgment, consumer); } catch (RuntimeException e) { Exception exception = new ConversionException("Failed to convert to message for: " + record, e); @@ -326,10 +327,11 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo } @Override - public void onMessage(List> records, Acknowledgment acknowledgment) { - Message message = null; + public void onMessage(List> records, Acknowledgment acknowledgment, + Consumer consumer) { + Message message = null; try { - message = toMessagingMessage(records, acknowledgment); + message = toMessagingMessage(records, acknowledgment, consumer); } catch (RuntimeException e) { Exception exception = new ConversionException("Failed to convert to message for: " + records, e); 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 54ef6de92e..fdee1105bd 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 @@ -33,9 +33,9 @@ import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAd 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.FilteringMessageListenerAdapter; import org.springframework.kafka.listener.adapter.RecordFilterStrategy; -import org.springframework.kafka.listener.adapter.RetryingAcknowledgingMessageListenerAdapter; +import org.springframework.kafka.listener.adapter.RetryingMessageListenerAdapter; import org.springframework.kafka.listener.config.ContainerProperties; import org.springframework.retry.support.RetryTemplate; import org.springframework.test.context.ContextConfiguration; @@ -112,7 +112,7 @@ public class KafkaMessageDrivenChannelAdapterParserTests { containerProps = TestUtils.getPropertyValue(container, "containerProperties", ContainerProperties.class); Object messageListener = containerProps.getMessageListener(); - assertThat(messageListener).isInstanceOf(FilteringAcknowledgingMessageListenerAdapter.class); + assertThat(messageListener).isInstanceOf(FilteringMessageListenerAdapter.class); Object delegate = TestUtils.getPropertyValue(messageListener, "delegate"); @@ -123,7 +123,7 @@ public class KafkaMessageDrivenChannelAdapterParserTests { adapter.afterPropertiesSet(); messageListener = containerProps.getMessageListener(); - assertThat(messageListener).isInstanceOf(RetryingAcknowledgingMessageListenerAdapter.class); + assertThat(messageListener).isInstanceOf(RetryingMessageListenerAdapter.class); delegate = TestUtils.getPropertyValue(messageListener, "delegate"); @@ -133,21 +133,21 @@ public class KafkaMessageDrivenChannelAdapterParserTests { adapter.afterPropertiesSet(); messageListener = containerProps.getMessageListener(); - assertThat(messageListener).isInstanceOf(FilteringAcknowledgingMessageListenerAdapter.class); + assertThat(messageListener).isInstanceOf(FilteringMessageListenerAdapter.class); delegate = TestUtils.getPropertyValue(messageListener, "delegate"); - assertThat(delegate).isInstanceOf(RetryingAcknowledgingMessageListenerAdapter.class); + assertThat(delegate).isInstanceOf(RetryingMessageListenerAdapter.class); adapter.setFilterInRetry(true); adapter.afterPropertiesSet(); messageListener = containerProps.getMessageListener(); - assertThat(messageListener).isInstanceOf(RetryingAcknowledgingMessageListenerAdapter.class); + assertThat(messageListener).isInstanceOf(RetryingMessageListenerAdapter.class); delegate = TestUtils.getPropertyValue(messageListener, "delegate"); - assertThat(delegate).isInstanceOf(FilteringAcknowledgingMessageListenerAdapter.class); + assertThat(delegate).isInstanceOf(FilteringMessageListenerAdapter.class); } } 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 c24fd42424..b1046254a4 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 @@ -23,6 +23,7 @@ import java.util.Arrays; import java.util.List; import java.util.Map; +import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.producer.ProducerRecord; @@ -89,8 +90,9 @@ public class MessageDrivenAdapterTests { adapter.setRecordMessageConverter(new MessagingMessageConverter() { @Override - public Message toMessage(ConsumerRecord record, Acknowledgment acknowledgment, Type type) { - Message message = super.toMessage(record, acknowledgment, type); + public Message toMessage(ConsumerRecord record, Acknowledgment acknowledgment, + Consumer consumer, Type type) { + Message message = super.toMessage(record, acknowledgment, consumer, type); return MessageBuilder.fromMessage(message).setHeader("testHeader", "testValue").build(); } @@ -136,7 +138,8 @@ public class MessageDrivenAdapterTests { adapter.setMessageConverter(new RecordMessageConverter() { @Override - public Message toMessage(ConsumerRecord record, Acknowledgment acknowledgment, Type payloadType) { + public Message toMessage(ConsumerRecord record, Acknowledgment acknowledgment, + Consumer consumer, Type type) { throw new RuntimeException("testError"); } @@ -173,8 +176,9 @@ public class MessageDrivenAdapterTests { adapter.setBatchMessageConverter(new BatchMessagingMessageConverter() { @Override - public Message toMessage(List> records, Acknowledgment acknowledgment, Type type) { - Message message = super.toMessage(records, acknowledgment, type); + public Message toMessage(List> records, Acknowledgment acknowledgment, + Consumer consumer, Type type) { + Message message = super.toMessage(records, acknowledgment, consumer, type); return MessageBuilder.fromMessage(message).setHeader("testHeader", "testValue").build(); } @@ -211,7 +215,7 @@ public class MessageDrivenAdapterTests { @Override public Message toMessage(List> records, Acknowledgment acknowledgment, - Type payloadType) { + Consumer consumer, Type payloadType) { throw new RuntimeException("testError"); }