@ -19,6 +19,7 @@ package org.springframework.boot.autoconfigure.kafka;
import org.springframework.boot.autoconfigure.kafka.KafkaProperties.Listener ;
import org.springframework.boot.autoconfigure.kafka.KafkaProperties.Listener ;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory ;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory ;
import org.springframework.kafka.core.ConsumerFactory ;
import org.springframework.kafka.core.ConsumerFactory ;
import org.springframework.kafka.core.KafkaTemplate ;
import org.springframework.kafka.listener.config.ContainerProperties ;
import org.springframework.kafka.listener.config.ContainerProperties ;
import org.springframework.kafka.support.converter.RecordMessageConverter ;
import org.springframework.kafka.support.converter.RecordMessageConverter ;
@ -35,6 +36,8 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer {
private RecordMessageConverter messageConverter ;
private RecordMessageConverter messageConverter ;
private KafkaTemplate < Object , Object > replyTemplate ;
/ * *
/ * *
* Set the { @link KafkaProperties } to use .
* Set the { @link KafkaProperties } to use .
* @param properties the properties
* @param properties the properties
@ -51,6 +54,14 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer {
this . messageConverter = messageConverter ;
this . messageConverter = messageConverter ;
}
}
/ * *
* Set the { @link KafkaTemplate } to use to send replies .
* @param replyTemplate the reply template
* /
void setReplyTemplate ( KafkaTemplate < Object , Object > replyTemplate ) {
this . replyTemplate = replyTemplate ;
}
/ * *
/ * *
* Configure the specified Kafka listener container factory . The factory can be
* Configure the specified Kafka listener container factory . The factory can be
* further tuned and default settings can be overridden .
* further tuned and default settings can be overridden .
@ -65,6 +76,9 @@ public class ConcurrentKafkaListenerContainerFactoryConfigurer {
if ( this . messageConverter ! = null ) {
if ( this . messageConverter ! = null ) {
listenerContainerFactory . setMessageConverter ( this . messageConverter ) ;
listenerContainerFactory . setMessageConverter ( this . messageConverter ) ;
}
}
if ( this . replyTemplate ! = null ) {
listenerContainerFactory . setReplyTemplate ( this . replyTemplate ) ;
}
Listener container = this . properties . getListener ( ) ;
Listener container = this . properties . getListener ( ) ;
ContainerProperties containerProperties = listenerContainerFactory
ContainerProperties containerProperties = listenerContainerFactory
. getContainerProperties ( ) ;
. getContainerProperties ( ) ;