From 037f9f440f219264105f1dd1dfcef41abb42c18d Mon Sep 17 00:00:00 2001 From: Fernando Blanch Calvete Date: Thu, 17 Sep 2026 10:57:22 +0200 Subject: [PATCH] fix: restore Kafka observation registry after bindings query Restore the original Kafka observation registry after serializing container properties for the bindings actuator endpoint. Add regression tests covering observation registry restoration and consumer restart after a bindings query. Fixes gh-3268 Signed-off-by: Fernando Blanch Calvete --- .../kafka/KafkaMessageChannelBinder.java | 18 ++++-- .../binder/kafka/KafkaConfigurationTests.java | 55 ++++++++++++++++++- 2 files changed, 64 insertions(+), 9 deletions(-) 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() {