diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java index 94665173f..87b1408ee 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java @@ -130,9 +130,10 @@ public class KafkaTopicProvisioner implements * {@link AdminClient}. */ public KafkaTopicProvisioner( - KafkaBinderConfigurationProperties kafkaBinderConfigurationProperties, - KafkaProperties kafkaProperties, - List adminClientConfigCustomizers) { + KafkaBinderConfigurationProperties kafkaBinderConfigurationProperties, + KafkaProperties kafkaProperties, + List adminClientConfigCustomizers) { + Assert.isTrue(kafkaProperties != null, "KafkaProperties cannot be null"); this.configurationProperties = kafkaBinderConfigurationProperties; this.adminClientProperties = kafkaProperties.buildAdminProperties(); @@ -143,6 +144,15 @@ public class KafkaTopicProvisioner implements adminClientConfigCustomizers.forEach(customizer -> customizer.configure(this.adminClientProperties)); } + /** + * Return an unmodifiable map of merged admin properties. + * @return the properties. + * @since 4.0.3 + */ + public Map getAdminClientProperties() { + return Collections.unmodifiableMap(this.adminClientProperties); + } + /** * Mutator for metadata retry operations. * @param metadataRetryOperations the retry configuration diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index f7ced7675..b7ce90faf 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -24,6 +24,7 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; import java.util.Collections; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Objects; @@ -98,6 +99,7 @@ import org.springframework.integration.support.MessageBuilder; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; +import org.springframework.kafka.core.KafkaAdmin; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.listener.AbstractMessageListenerContainer; @@ -237,6 +239,8 @@ public class KafkaMessageChannelBinder extends private final List> kafkaMessageListenerContainers = new ArrayList<>(); + private final KafkaAdmin kafkaAdmin; + public KafkaMessageChannelBinder( KafkaBinderConfigurationProperties configurationProperties, KafkaTopicProvisioner provisioningProvider) { @@ -281,6 +285,7 @@ public class KafkaMessageChannelBinder extends this.rebalanceListener = rebalanceListener; this.dlqPartitionFunction = dlqPartitionFunction; this.dlqDestinationResolver = dlqDestinationResolver; + this.kafkaAdmin = new KafkaAdmin(new HashMap<>(provisioningProvider.getAdminClientProperties())); } private static String[] headersToMap( @@ -503,9 +508,8 @@ public class KafkaMessageChannelBinder extends } handler.setHeaderMapper(mapper); - if (this.configurationProperties.isEnableObservation()) { - kafkaTemplate.setObservationEnabled(true); - } + kafkaTemplate.setObservationEnabled(this.configurationProperties.isEnableObservation()); + kafkaTemplate.setKafkaAdmin(this.kafkaAdmin); kafkaTemplate.setApplicationContext(getApplicationContext()); return handler; @@ -632,9 +636,7 @@ public class KafkaMessageChannelBinder extends : new ContainerProperties(topics) : new ContainerProperties(topicPartitionOffsets); - if (this.configurationProperties.isEnableObservation()) { - containerProperties.setObservationEnabled(true); - } + containerProperties.setObservationEnabled(this.configurationProperties.isEnableObservation()); KafkaAwareTransactionManager transMan = transactionManager( extendedConsumerProperties.getExtension().getTransactionManager()); @@ -668,6 +670,7 @@ public class KafkaMessageChannelBinder extends }; this.kafkaMessageListenerContainers.add(messageListenerContainer); + messageListenerContainer.setKafkaAdmin(this.kafkaAdmin); messageListenerContainer.setConcurrency(concurrency); // these won't be needed if the container is made a bean AbstractApplicationContext applicationContext = getApplicationContext(); @@ -1126,14 +1129,16 @@ public class KafkaMessageChannelBinder extends .getDlqProducerProperties(); KafkaAwareTransactionManager transMan = transactionManager( properties.getExtension().getTransactionManager()); - final ExtendedProducerProperties producerProperties = new ExtendedProducerProperties<>(dlqProducerProperties); + final ExtendedProducerProperties producerProperties = + new ExtendedProducerProperties<>(dlqProducerProperties); producerProperties.populateBindingName(properties.getBindingName()); ProducerFactory producerFactory = transMan != null ? transMan.getProducerFactory() : getProducerFactory(null, producerProperties, destination.getName() + ".dlq.producer", destination.getName()); - final KafkaTemplate kafkaTemplate = new KafkaTemplate<>( - producerFactory); + final KafkaTemplate kafkaTemplate = new KafkaTemplate<>(producerFactory); + kafkaTemplate.setObservationEnabled(this.configurationProperties.isEnableObservation()); + kafkaTemplate.setKafkaAdmin(this.kafkaAdmin); Object timeout = producerFactory.getConfigurationProperties().get(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG); Long sendTimeout = null; diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index 0bc04df02..bf5b6575a 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -119,6 +119,7 @@ import org.springframework.integration.kafka.support.KafkaSendFailureException; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; +import org.springframework.kafka.core.KafkaAdmin; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.listener.AbstractMessageListenerContainer; @@ -156,6 +157,7 @@ import org.springframework.util.backoff.FixedBackOff; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatExceptionOfType; +import static org.assertj.core.api.Assertions.entry; import static org.assertj.core.api.Assertions.fail; import static org.mockito.Mockito.mock; @@ -346,6 +348,40 @@ public class KafkaBinderTests extends return new DefaultKafkaConsumerFactory<>(props, keyDecoder, valueDecoder); } + @SuppressWarnings({ "rawtypes", "unchecked" }) + @Test + void bindersAdmin() throws Exception { + KafkaBinderConfigurationProperties props = createConfigurationProperties(); + props.getConfiguration().put(AdminClientConfig.CLIENT_ID_CONFIG, "binder"); + props.setEnableObservation(true); + Binder binder = getBinder(props); + BindingProperties producerBindingProperties = createProducerBindingProperties( + createProducerProperties()); + DirectChannel moduleOutputChannel = createBindableChannel("output", + producerBindingProperties); + + ExtendedConsumerProperties consumerProperties = createConsumerProperties(); + DirectChannel moduleInputChannel = createBindableChannel("input", + createConsumerBindingProperties(consumerProperties)); + + Binding producerBinding = binder.bindProducer("admin.0", + moduleOutputChannel, producerBindingProperties.getProducer()); + + Binding consumerBinding = binder.bindConsumer("admin.0", + "testSendAndReceiveNoOriginalContentType", moduleInputChannel, + consumerProperties); + + assertThat( + KafkaTestUtils.getPropertyValue(producerBinding, "lifecycle.kafkaTemplate.kafkaAdmin", KafkaAdmin.class) + .getConfigurationProperties()).contains(entry(AdminClientConfig.CLIENT_ID_CONFIG, "binder")); + assertThat(KafkaTestUtils + .getPropertyValue(consumerBinding, "lifecycle.messageListenerContainer.kafkaAdmin", KafkaAdmin.class) + .getConfigurationProperties()).contains(entry(AdminClientConfig.CLIENT_ID_CONFIG, "binder")); + + consumerBinding.unbind(); + producerBinding.unbind(); + } + @SuppressWarnings({ "rawtypes", "unchecked" }) @Test void testDefaultHeaderMapper() throws Exception {