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 b7c342a73f..625f8a5f0f 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 @@ -23,11 +23,12 @@ 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.RecordMessagingMessageListenerAdapter; import org.springframework.kafka.listener.adapter.RetryingAcknowledgingMessageListenerAdapter; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.converter.MessageConverter; +import org.springframework.kafka.support.converter.RecordMessageConverter; import org.springframework.messaging.Message; import org.springframework.retry.RecoveryCallback; import org.springframework.retry.support.RetryTemplate; @@ -48,7 +49,7 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo private final AbstractMessageListenerContainer messageListenerContainer; - private final MessagingMessageListenerAdapter listener = new IntegrationMessageListener(); + private final RecordMessagingMessageListenerAdapter listener = new IntegrationMessageListener(); private RecordFilterStrategy recordFilterStrategy; @@ -68,7 +69,23 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo this.messageListenerContainer.setAutoStartup(false); } + /** + * Set the message converter; must be a {@link RecordMessageConverter}. + * @param messageConverter the converter. + * @deprecated in favor of {@link #setRecordMessageConverter(RecordMessageConverter)}. + */ + @Deprecated public void setMessageConverter(MessageConverter messageConverter) { + Assert.isInstanceOf(RecordMessageConverter.class, messageConverter); + this.listener.setMessageConverter((RecordMessageConverter) messageConverter); + } + + /** + * Set the message converter to use with a record-based consumer. + * @param messageConverter the converter. + * @since 2.1 + */ + public void setRecordMessageConverter(RecordMessageConverter messageConverter) { this.listener.setMessageConverter(messageConverter); } @@ -185,10 +202,10 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo return getPhase(); } - private class IntegrationMessageListener extends MessagingMessageListenerAdapter { + private class IntegrationMessageListener extends RecordMessagingMessageListenerAdapter { IntegrationMessageListener() { - super(null); + super(null, null); } @Override 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 b6e54c48a1..a04f6cae91 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 @@ -68,7 +68,7 @@ public class MessageDrivenAdapterTests { QueueChannel out = new QueueChannel(); adapter.setOutputChannel(out); adapter.afterPropertiesSet(); - adapter.setMessageConverter(new MessagingMessageConverter() { + adapter.setRecordMessageConverter(new MessagingMessageConverter() { @Override public Message toMessage(ConsumerRecord record, Acknowledgment acknowledgment, Type type) {