diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/ObservationAutoConfiguration.java b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/ObservationAutoConfiguration.java new file mode 100644 index 000000000..25e3158dc --- /dev/null +++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/ObservationAutoConfiguration.java @@ -0,0 +1,60 @@ +/* + * Copyright 2022-2022 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.stream.binder.rabbit.config; + +import org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer; +import org.springframework.amqp.rabbit.listener.MessageListenerContainer; +import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; +import org.springframework.cloud.stream.config.ListenerContainerCustomizer; +import org.springframework.cloud.stream.config.ProducerMessageHandlerCustomizer; +import org.springframework.context.ApplicationContext; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.core.Ordered; +import org.springframework.core.annotation.Order; +import org.springframework.integration.amqp.outbound.AmqpOutboundEndpoint; +import org.springframework.messaging.MessageHandler; + +@Configuration(proxyBeanMethods = false) +@ConditionalOnBean(org.springframework.boot.actuate.autoconfigure.observation.ObservationAutoConfiguration.class) +public class ObservationAutoConfiguration { + + @Bean + @Order(Ordered.HIGHEST_PRECEDENCE) + ListenerContainerCustomizer observedListenerContainerCustomizer( + ApplicationContext applicationContext) { + return (container, destinationName, group) -> { + if (container instanceof AbstractMessageListenerContainer) { + AbstractMessageListenerContainer abstractMessageListenerContainer = ((AbstractMessageListenerContainer) container); + abstractMessageListenerContainer.setObservationEnabled(true); + abstractMessageListenerContainer.setApplicationContext(applicationContext); + } + }; + } + + @Bean + @Order(Ordered.HIGHEST_PRECEDENCE) + ProducerMessageHandlerCustomizer observedProducerMessageHandlerCustomizer( + ApplicationContext applicationContext) { + return (handler, destinationName) -> { + if (handler instanceof AmqpOutboundEndpoint) { + ((AmqpOutboundEndpoint) handler).getRabbitTemplate().setObservationEnabled(true); + ((AmqpOutboundEndpoint) handler).getRabbitTemplate().setApplicationContext(applicationContext); + } + }; + } +} diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitMessageChannelBinderConfiguration.java b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitMessageChannelBinderConfiguration.java index 62ad835fe..c43eeb2fd 100644 --- a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitMessageChannelBinderConfiguration.java +++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/config/RabbitMessageChannelBinderConfiguration.java @@ -44,8 +44,9 @@ import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; import org.springframework.integration.amqp.inbound.AmqpInboundChannelAdapter; import org.springframework.integration.amqp.inbound.AmqpMessageSource; -import org.springframework.integration.amqp.outbound.AmqpOutboundEndpoint; import org.springframework.lang.Nullable; +import org.springframework.messaging.MessageHandler; +import org.springframework.util.CollectionUtils; /** @@ -78,9 +79,9 @@ public class RabbitMessageChannelBinderConfiguration { @Bean RabbitMessageChannelBinder rabbitMessageChannelBinder( - @Nullable ListenerContainerCustomizer listenerContainerCustomizer, + @Nullable List> listenerContainerCustomizers, @Nullable MessageSourceCustomizer sourceCustomizer, - @Nullable ProducerMessageHandlerCustomizer producerMessageHandlerCustomizer, + @Nullable List> producerMessageHandlerCustomizers, @Nullable ConsumerEndpointCustomizer consumerCustomizer, List declarableCustomizers, @Nullable ConnectionNameStrategy connectionNameStrategy) { @@ -91,9 +92,31 @@ public class RabbitMessageChannelBinderConfiguration { ((AbstractConnectionFactory) this.rabbitConnectionFactory).setConnectionNameStrategy(f -> connectionNamePrefix + "#" + nameIncrementer.getAndIncrement()); } + + ListenerContainerCustomizer composedCistomizer = new ListenerContainerCustomizer<>() { + @Override + public void configure(MessageListenerContainer container, String destinationName, String group) { + if (!CollectionUtils.isEmpty(listenerContainerCustomizers)) { + for (ListenerContainerCustomizer customizer : listenerContainerCustomizers) { + customizer.configure(container, destinationName, group); + } + } + } + }; + ProducerMessageHandlerCustomizer producerMessageHandlerCustomizer = new ProducerMessageHandlerCustomizer<>() { + @Override + public void configure(MessageHandler handler, String destinationName) { + if (!CollectionUtils.isEmpty(producerMessageHandlerCustomizers)) { + for (ProducerMessageHandlerCustomizer customizer : producerMessageHandlerCustomizers) { + customizer.configure(handler, destinationName); + } + } + } + }; + RabbitMessageChannelBinder binder = new RabbitMessageChannelBinder( this.rabbitConnectionFactory, this.rabbitProperties, - provisioningProvider(declarableCustomizers), listenerContainerCustomizer, sourceCustomizer); + provisioningProvider(declarableCustomizers), composedCistomizer, sourceCustomizer); binder.setAdminAddresses( this.rabbitBinderConfigurationProperties.getAdminAddresses()); binder.setCompressingPostProcessor(gZipPostProcessor()); diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports index a6aa77f21..293619926 100644 --- a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports +++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports @@ -1 +1,2 @@ org.springframework.cloud.stream.binder.rabbit.config.ExtendedBindingHandlerMappingsProviderConfiguration +org.springframework.cloud.stream.binder.rabbit.config.ObservationAutoConfiguration diff --git a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/integration/RabbitBinderModuleTests.java b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/integration/RabbitBinderModuleTests.java index 52a91fe40..451c91a48 100644 --- a/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/integration/RabbitBinderModuleTests.java +++ b/binders/rabbit-binder/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/integration/RabbitBinderModuleTests.java @@ -71,9 +71,10 @@ import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.integration.amqp.inbound.AmqpInboundChannelAdapter; import org.springframework.integration.amqp.inbound.AmqpMessageSource; -import org.springframework.integration.amqp.outbound.AmqpOutboundEndpoint; import org.springframework.integration.channel.DirectChannel; +import org.springframework.integration.handler.AbstractMessageHandler; import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.MessageHandler; import org.springframework.messaging.support.GenericMessage; import org.springframework.retry.backoff.ExponentialBackOffPolicy; import org.springframework.retry.policy.SimpleRetryPolicy; @@ -421,8 +422,10 @@ public class RabbitBinderModuleTests { @Bean public ListenerContainerCustomizer containerCustomizer() { - return (c, q, g) -> ((AbstractMessageListenerContainer) c).setBeanName( - "setByCustomizerForQueue:" + q + (g == null ? "" : ",andGroup:" + g)); + return (c, q, g) -> { + ((AbstractMessageListenerContainer) c).setBeanName( + "setByCustomizerForQueue:" + q + (g == null ? "" : ",andGroup:" + g)); + }; } @Bean @@ -431,8 +434,12 @@ public class RabbitBinderModuleTests { } @Bean - public ProducerMessageHandlerCustomizer messageHandlerCustomizer() { - return (handler, destinationName) -> handler.setBeanName("setByCustomizer:" + destinationName); + public ProducerMessageHandlerCustomizer messageHandlerCustomizer() { + return (handler, destinationName) -> { + if (handler instanceof AbstractMessageHandler) { + ((AbstractMessageHandler) handler).setBeanName("setByCustomizer:" + destinationName); + } + }; } @Bean diff --git a/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index 07da13232..524c95614 100644 --- a/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/core/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -1475,8 +1475,7 @@ public class ImplicitFunctionBindingTests { return new FunctionAroundWrapper() { @Override - protected Object doApply(Object input, - FunctionInvocationWrapper targetFunction) { + protected Object doApply(Object input, FunctionInvocationWrapper targetFunction) { return targetFunction.apply(input); } };