GH-2673: Use Binder Admin Config with Observation

Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/2673

Also configure observation (if enabled) on DLQ template.
Resolves #2703
This commit is contained in:
Gary Russell
2023-04-10 16:52:06 -04:00
committed by Oleg Zhurakousky
parent dd8e707e01
commit 57733739e5
3 changed files with 63 additions and 12 deletions

View File

@@ -130,9 +130,10 @@ public class KafkaTopicProvisioner implements
* {@link AdminClient}.
*/
public KafkaTopicProvisioner(
KafkaBinderConfigurationProperties kafkaBinderConfigurationProperties,
KafkaProperties kafkaProperties,
List<AdminClientConfigCustomizer> adminClientConfigCustomizers) {
KafkaBinderConfigurationProperties kafkaBinderConfigurationProperties,
KafkaProperties kafkaProperties,
List<AdminClientConfigCustomizer> 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<String, Object> getAdminClientProperties() {
return Collections.unmodifiableMap(this.adminClientProperties);
}
/**
* Mutator for metadata retry operations.
* @param metadataRetryOperations the retry configuration

View File

@@ -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<AbstractMessageListenerContainer<?, ?>> 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<byte[], byte[]> 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<byte[], byte[]> transMan = transactionManager(
properties.getExtension().getTransactionManager());
final ExtendedProducerProperties<KafkaProducerProperties> producerProperties = new ExtendedProducerProperties<>(dlqProducerProperties);
final ExtendedProducerProperties<KafkaProducerProperties> 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;

View File

@@ -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<KafkaConsumerProperties> consumerProperties = createConsumerProperties();
DirectChannel moduleInputChannel = createBindableChannel("input",
createConsumerBindingProperties(consumerProperties));
Binding<MessageChannel> producerBinding = binder.bindProducer("admin.0",
moduleOutputChannel, producerBindingProperties.getProducer());
Binding<MessageChannel> 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 {