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 c54ed08796..9d6a95b047 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 @@ -914,12 +914,18 @@ public ObservationConfig observationConfig() { @Override protected Map doGetAdditionalConfigurationProperties(String destinationName) { ContainerProperties kafkaContainerProperties = this.kafkaMessageListenerContainers.iterator().next().getContainerProperties(); - // see 3167 we need to nullify ObservationRegistry to avoid jackson deserialization error - kafkaContainerProperties.setObservationRegistry(new NullObservationRegistry()); - Map mapOfContainerProperties = this.objectMapper.convertValue(kafkaContainerProperties, Map.class); - Map additionalConfigurationProperties = new HashMap<>(); - additionalConfigurationProperties.put("containerProperties", mapOfContainerProperties); - return additionalConfigurationProperties; + ObservationRegistry observationRegistry = kafkaContainerProperties.getObservationRegistry(); + try { + // See 3167; nullify ObservationRegistry to avoid Jackson deserialization errors. + kafkaContainerProperties.setObservationRegistry(new NullObservationRegistry()); + Map mapOfContainerProperties = this.objectMapper.convertValue(kafkaContainerProperties, Map.class); + Map additionalConfigurationProperties = new HashMap<>(); + additionalConfigurationProperties.put("containerProperties", mapOfContainerProperties); + return additionalConfigurationProperties; + } + finally { + kafkaContainerProperties.setObservationRegistry(observationRegistry); + } } private BiFunction, Exception, TopicPartition> createDestResolver( diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaConfigurationTests.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaConfigurationTests.java index 2a02415657..a1b9eeae41 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaConfigurationTests.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaConfigurationTests.java @@ -16,18 +16,25 @@ package org.springframework.cloud.stream.binder.kafka; +import java.time.Duration; import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Consumer; +import io.micrometer.observation.ObservationRegistry; +import org.awaitility.Awaitility; import org.junit.jupiter.api.Test; +import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.binder.Binding; import org.springframework.cloud.stream.binding.BindingService; +import org.springframework.cloud.stream.endpoint.BindingsEndpoint; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.test.context.EmbeddedKafka; import org.springframework.test.annotation.DirtiesContext; @@ -35,20 +42,33 @@ /** * @author Oleg Zhurakousky + * @author Fernando Blanch */ @SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.NONE, properties = { "spring.cloud.function.definition=barConsumer;fooConsumer", "spring.kafka.listener.immediate-stop=true", - "spring.cloud.stream.bindings.fooConsumer-in-0.destination=foo" + "spring.cloud.stream.kafka.binder.enableObservation=true", + "management.endpoint.bindings.enabled=true", + "management.endpoints.web.exposure.include=bindings", + "spring.cloud.stream.bindings.fooConsumer-in-0.destination=foo", + "spring.cloud.stream.bindings.barConsumer-in-0.destination=bar", + "spring.cloud.stream.bindings.barConsumer-in-0.group=bar-group" }) @EmbeddedKafka @DirtiesContext public class KafkaConfigurationTests { + private static final AtomicInteger RECEIVED_MESSAGES = new AtomicInteger(); @Autowired private BindingService bindingService; + @Autowired + private KafkaTemplate kafkaTemplate; + + @Autowired + private BindingsEndpoint bindingsEndpoint; + @Test void testKafkaContainerConfigurationPropagation() throws Exception { Binding fooDestination = this.bindingService.getConsumerBindings("fooConsumer-in-0").iterator().next(); @@ -60,14 +80,43 @@ void testKafkaContainerConfigurationPropagation() throws Exception { assertThat(((Map) barAdditionalConfigurationProperties.get("containerProperties")).get("stopImmediate")).isEqualTo(true); } + @Test + void testObservationRegistryIsRestoredAfterGettingAdditionalConfigurationProperties() { + Binding binding = this.bindingService.getConsumerBindings("barConsumer-in-0").iterator().next(); + DirectFieldAccessor bindingAccessor = new DirectFieldAccessor(binding); + ObservationRegistry observationRegistry = (ObservationRegistry) bindingAccessor + .getPropertyValue("lifecycle.messageListenerContainer.containerProperties.observationRegistry"); + + this.bindingsEndpoint.queryStates(); + + assertThat(bindingAccessor.getPropertyValue( + "lifecycle.messageListenerContainer.containerProperties.observationRegistry")) + .isSameAs(observationRegistry); + } + + @Test + void testKafkaConsumerCanBeRestartedAfterGettingAdditionalConfigurationProperties() { + RECEIVED_MESSAGES.set(0); + Binding binding = this.bindingService.getConsumerBindings("barConsumer-in-0").iterator().next(); + + this.bindingsEndpoint.queryStates(); + binding.stop(); + binding.start(); + + this.kafkaTemplate.send("bar", null, "foo".getBytes()); + this.kafkaTemplate.flush(); + + Awaitility.await().atMost(Duration.ofSeconds(10)) + .untilAsserted(() -> assertThat(RECEIVED_MESSAGES).hasValue(1)); + } + @EnableAutoConfiguration @Configuration public static class Config { @Bean Consumer barConsumer() { - return message -> { - }; + return message -> RECEIVED_MESSAGES.incrementAndGet(); } @Bean Consumer fooConsumer() {