diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java index df197ce4..6de21f52 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java @@ -364,7 +364,7 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener private final String consumerGroupId = this.containerProperties.getGroupId() == null ? (String) KafkaMessageListenerContainer.this.consumerFactory.getConfigurationProperties() - .get(ConsumerConfig.GROUP_ID_CONFIG) + .get(ConsumerConfig.GROUP_ID_CONFIG) : this.containerProperties.getGroupId(); private final TaskScheduler taskScheduler; @@ -463,6 +463,9 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener } this.monitorTask = this.taskScheduler.scheduleAtFixedRate(() -> checkConsumer(), this.containerProperties.getMonitorInterval() * 1000); + if (this.containerProperties.isLogContainerConfig() && this.logger.isInfoEnabled()) { + this.logger.info(this); + } } protected void checkConsumer() { @@ -1260,6 +1263,19 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener this.seeks.add(new TopicPartitionInitialOffset(topic, partition, SeekPosition.END)); } + @Override + public String toString() { + return "KafkaMessageListenerContainer.ListenerConsumer [" + + "containerProperties=" + this.containerProperties + + ", listenerType=" + this.listenerType + + ", isConsumerAwareListener=" + this.isConsumerAwareListener + + ", isBatchListener=" + this.isBatchListener + + ", autoCommit=" + this.autoCommit + + ", consumerGroupId=" + this.consumerGroupId + + ", clientIdSuffix=" + KafkaMessageListenerContainer.this.clientIdSuffix + + "]"; + } + private final class ConsumerAcknowledgment implements Acknowledgment { private final ConsumerRecord record; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/config/ContainerProperties.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/config/ContainerProperties.java index 93ef26e7..1d9362df 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/config/ContainerProperties.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/config/ContainerProperties.java @@ -33,6 +33,7 @@ import org.springframework.kafka.support.TopicPartitionInitialOffset; import org.springframework.scheduling.TaskScheduler; import org.springframework.transaction.PlatformTransactionManager; import org.springframework.util.Assert; +import org.springframework.util.StringUtils; /** * Contains runtime properties for a listener container. @@ -159,6 +160,8 @@ public class ContainerProperties { private String clientId = ""; + private boolean logContainerConfig; + public ContainerProperties(String... topics) { Assert.notEmpty(topics, "An array of topicPartitions must be provided"); this.topics = Arrays.asList(topics).toArray(new String[topics.length]); @@ -204,6 +207,7 @@ public class ContainerProperties { * @param ackMode the {@link AckMode}; default BATCH. */ public void setAckMode(AbstractMessageListenerContainer.AckMode ackMode) { + Assert.notNull(ackMode, "'ackMode' cannot be null"); this.ackMode = ackMode; } @@ -483,4 +487,55 @@ public class ContainerProperties { this.clientId = clientId; } + /** + * Log the container configuration if true (INFO). + * @return true to log. + * @since 2.0.1 + */ + public boolean isLogContainerConfig() { + return this.logContainerConfig; + } + + /** + * Set to true to instruct each container to log this configuration. + * @param logContainerConfig true to log. + * @since 2.1.1 + */ + public void setLogContainerConfig(boolean logContainerConfig) { + this.logContainerConfig = logContainerConfig; + } + + @Override + public String toString() { + return "ContainerProperties [" + + (this.topics != null ? "topics=" + Arrays.toString(this.topics) : "") + + (this.topicPattern != null ? ", topicPattern=" + this.topicPattern : "") + + (this.topicPartitions != null + ? ", topicPartitions=" + Arrays.toString(this.topicPartitions) : "") + + ", ackMode=" + this.ackMode + + ", ackCount=" + this.ackCount + + ", ackTime=" + this.ackTime + + ", messageListener=" + this.messageListener + + ", pollTimeout=" + this.pollTimeout + + (this.consumerTaskExecutor != null + ? ", consumerTaskExecutor=" + this.consumerTaskExecutor : "") + + (this.errorHandler != null ? ", errorHandler=" + this.errorHandler : "") + + ", shutdownTimeout=" + this.shutdownTimeout + + (this.consumerRebalanceListener != null + ? ", consumerRebalanceListener=" + this.consumerRebalanceListener : "") + + (this.commitCallback != null ? ", commitCallback=" + this.commitCallback : "") + + ", syncCommits=" + this.syncCommits + + ", ackOnError=" + this.ackOnError + + ", idleEventInterval=" + + (this.idleEventInterval == null ? "not enabled" : this.idleEventInterval) + + (this.groupId != null ? ", groupId=" + this.groupId : "") + + (this.transactionManager != null + ? ", transactionManager=" + this.transactionManager : "") + + ", monitorInterval=" + this.monitorInterval + + (this.scheduler != null ? ", scheduler=" + this.scheduler : "") + + ", noPollThreshold=" + this.noPollThreshold + + (StringUtils.hasText(this.clientId) ? ", clientId=" + this.clientId : "") + + "]"; + } + } diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java index f4a375a5..acc081ff 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java @@ -102,6 +102,7 @@ public class ConcurrentMessageListenerContainerTests { Map props = KafkaTestUtils.consumerProps("test1", "true", embeddedKafka); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(props); ContainerProperties containerProps = new ContainerProperties(topic1); + containerProps.setLogContainerConfig(true); final CountDownLatch latch = new CountDownLatch(4); final Set listenerThreadNames = new ConcurrentSkipListSet<>(); diff --git a/src/reference/asciidoc/kafka.adoc b/src/reference/asciidoc/kafka.adoc index 40146686..f2a521f5 100644 --- a/src/reference/asciidoc/kafka.adoc +++ b/src/reference/asciidoc/kafka.adoc @@ -496,6 +496,8 @@ return container; Refer to the JavaDocs for `ContainerProperties` for more information about the various properties that can be set. +Since version _2.1.1_, a new property `logContainerConfig` is available; when true, and INFO logging is enabled, each listener container will write a log message summarizing its configuration properties. + ====== ConcurrentMessageListenerContainer The single constructor is similar to the first `KafkaListenerContainer` constructor: