From 792b8fe6d11327b1616201f87d4d68b4192545ba Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 27 Feb 2018 16:04:54 -0500 Subject: [PATCH] INT-4416: Rework Micrometer Metrics JIRA: https://jira.spring.io/browse/INT-4416 Metrics should be under a common name, discriminated with tags. --- build.gradle | 2 +- .../channel/AbstractMessageChannel.java | 55 ++++++-- .../channel/AbstractPollableChannel.java | 20 +++ .../integration/channel/NullChannel.java | 21 ++++ .../endpoint/AbstractMessageSource.java | 22 +++- .../handler/AbstractMessageHandler.java | 52 +++++--- .../management/IntegrationManagement.java | 17 +++ .../IntegrationManagementConfigurer.java | 53 ++++++++ .../micrometer/MicrometerMetricsFactory.java | 3 + .../micrometer/MicrometerMetricsTests.java | 119 +++++++----------- 10 files changed, 256 insertions(+), 108 deletions(-) diff --git a/build.gradle b/build.gradle index 3ac5d84cd7..76ec595c27 100644 --- a/build.gradle +++ b/build.gradle @@ -120,7 +120,7 @@ subprojects { subproject -> jythonVersion = '2.5.3' kryoShadedVersion = '3.0.3' log4jVersion = '2.10.0' - micrometerVersion = '1.0.0' + micrometerVersion = '1.0.1' mockitoVersion = '2.11.0' mysqlVersion = '6.0.6' pahoMqttClientVersion = '1.2.0' diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java index 31b98b8db8..5ee328ec53 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java @@ -48,6 +48,10 @@ import org.springframework.messaging.support.ChannelInterceptor; import org.springframework.util.Assert; import org.springframework.util.StringUtils; +import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.core.instrument.Timer; +import io.micrometer.core.instrument.Timer.Sample; + /** * Base class for {@link MessageChannel} implementations providing common * properties such as the channel name. Also provides the common functionality @@ -86,6 +90,8 @@ public abstract class AbstractMessageChannel extends IntegrationObjectSupport private volatile AbstractMessageChannelMetrics channelMetrics = new DefaultMessageChannelMetrics(); + private MeterRegistry meterRegistry; + public AbstractMessageChannel() { this.interceptors = new ChannelInterceptorList(logger); } @@ -100,6 +106,15 @@ public abstract class AbstractMessageChannel extends IntegrationObjectSupport this.shouldTrack = shouldTrack; } + @Override + public void registerMeterRegistry(MeterRegistry registry) { + this.meterRegistry = registry; + } + + protected MeterRegistry getMeterRegistry() { + return this.meterRegistry; + } + @Override public void setCountsEnabled(boolean countsEnabled) { this.countsEnabled = countsEnabled; @@ -417,6 +432,10 @@ public abstract class AbstractMessageChannel extends IntegrationObjectSupport boolean countsEnabled = this.countsEnabled; ChannelInterceptorList interceptors = this.interceptors; AbstractMessageChannelMetrics channelMetrics = this.channelMetrics; + Sample sample = null; + if (this.meterRegistry != null) { + sample = Timer.start(this.meterRegistry); + } try { if (this.datatypes.length > 0) { message = this.convertPayloadIfNecessary(message); @@ -433,16 +452,22 @@ public abstract class AbstractMessageChannel extends IntegrationObjectSupport } } if (countsEnabled) { - if (channelMetrics.getTimer() != null) { - final Message messageToSend = message; - sent = channelMetrics.getTimer().recordCallable(() -> doSend(messageToSend, timeout)); + metrics = channelMetrics.beforeSend(); + if (this.meterRegistry != null) { + sample = Timer.start(this.meterRegistry); } - else { - metrics = channelMetrics.beforeSend(); - sent = doSend(message, timeout); - channelMetrics.afterSend(metrics, sent); - metricsProcessed = true; + sent = doSend(message, timeout); + if (sample != null) { + sample.stop(Timer.builder(SEND_TIMER_NAME) + .tag("type", "channel") + .tag("name", getComponentName() == null ? "unknown" : getComponentName()) + .tag("result", sent ? "success" : "failure") + .tag("exception", "none") + .description("Subflow process time") + .register(this.meterRegistry)); } + channelMetrics.afterSend(metrics, sent); + metricsProcessed = true; } else { sent = doSend(message, timeout); @@ -459,12 +484,16 @@ public abstract class AbstractMessageChannel extends IntegrationObjectSupport } catch (Exception e) { if (countsEnabled && !metricsProcessed) { - if (channelMetrics.getErrorCounter() != null) { - channelMetrics.getErrorCounter().increment(); - } - else { - channelMetrics.afterSend(metrics, false); + if (sample != null) { + sample.stop(Timer.builder(SEND_TIMER_NAME) + .tag("type", "channel") + .tag("name", getComponentName() == null ? "unknown" : getComponentName()) + .tag("result", "failure") + .tag("exception", e.getClass().getSimpleName()) + .description("Subflow process time") + .register(this.meterRegistry)); } + channelMetrics.afterSend(metrics, false); } if (interceptorStack != null) { interceptors.afterSendCompletion(message, this, sent, e, interceptorStack); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractPollableChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractPollableChannel.java index 1fd53704ca..586c2f3518 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractPollableChannel.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractPollableChannel.java @@ -27,6 +27,8 @@ import org.springframework.messaging.support.ChannelInterceptor; import org.springframework.messaging.support.ExecutorChannelInterceptor; import org.springframework.util.CollectionUtils; +import io.micrometer.core.instrument.Counter; + /** * Base class for all pollable channels. * @@ -104,6 +106,15 @@ public abstract class AbstractPollableChannel extends AbstractMessageChannel } Message message = this.doReceive(timeout); if (countsEnabled && message != null) { + if (getMeterRegistry() != null) { + Counter.builder(RECEIVE_COUNTER_NAME) + .tag("name", getComponentName()) + .tag("type", "channel") + .tag("result", "success") + .tag("exception", "none") + .description("Messages received") + .register(getMeterRegistry()).increment(); + } getMetrics().afterReceive(); counted = true; } @@ -121,6 +132,15 @@ public abstract class AbstractPollableChannel extends AbstractMessageChannel } catch (RuntimeException e) { if (countsEnabled && !counted) { + if (getMeterRegistry() != null) { + Counter.builder(RECEIVE_COUNTER_NAME) + .tag("name", getComponentName() == null ? "unknown" : getComponentName()) + .tag("type", "channel") + .tag("result", "failure") + .tag("exception", e.getClass().getSimpleName()) + .description("Messages received") + .register(getMeterRegistry()).increment(); + } getMetrics().afterError(); } if (!CollectionUtils.isEmpty(interceptorStack)) { diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/NullChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/NullChannel.java index d83e86138a..b7acbddafd 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/NullChannel.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/NullChannel.java @@ -16,6 +16,8 @@ package org.springframework.integration.channel; +import java.util.concurrent.TimeUnit; + import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -32,6 +34,9 @@ import org.springframework.messaging.PollableChannel; import org.springframework.util.Assert; import org.springframework.util.StringUtils; +import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.core.instrument.Timer; + /** * A channel implementation that essentially behaves like "/dev/null". * All receive() calls will return null, and all send() calls @@ -59,6 +64,8 @@ public class NullChannel implements PollableChannel, MessageChannelMetrics, private String beanName; + private MeterRegistry meterRegistry; + @Override public void setBeanName(String beanName) { this.beanName = beanName; @@ -86,6 +93,11 @@ public class NullChannel implements PollableChannel, MessageChannelMetrics, return "channel"; } + @Override + public void registerMeterRegistry(MeterRegistry registry) { + this.meterRegistry = registry; + } + @Override public void configureMetrics(AbstractMessageChannelMetrics metrics) { Assert.notNull(metrics, "'metrics' must not be null"); @@ -215,6 +227,15 @@ public class NullChannel implements PollableChannel, MessageChannelMetrics, this.logger.debug("message sent to null channel: " + message); } if (this.countsEnabled) { + if (this.meterRegistry != null) { + Timer.builder(SEND_TIMER_NAME) + .tag("type", "channel") + .tag("name", getComponentName() == null ? "unknown" : getComponentName()) + .tag("result", "success") + .tag("exception", "none") + .description("Subflow process time") + .register(this.meterRegistry).record(0, TimeUnit.MILLISECONDS); + } this.channelMetrics.afterSend(this.channelMetrics.beforeSend(), true); } return true; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractMessageSource.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractMessageSource.java index 48f619e632..e4af49eed1 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractMessageSource.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractMessageSource.java @@ -34,6 +34,7 @@ import org.springframework.messaging.MessagingException; import org.springframework.util.CollectionUtils; import io.micrometer.core.instrument.Counter; +import io.micrometer.core.instrument.MeterRegistry; /** * @author Mark Fisher @@ -65,11 +66,18 @@ public abstract class AbstractMessageSource extends AbstractExpressionEvaluat private volatile boolean loggingEnabled = true; + private MeterRegistry meterRegistry; + public void setHeaderExpressions(Map headerExpressions) { this.headerExpressions = (headerExpressions != null) ? headerExpressions : Collections.emptyMap(); } + @Override + public void registerMeterRegistry(MeterRegistry registry) { + this.meterRegistry = registry; + } + @Override public void setBeanName(String name) { this.beanName = name; @@ -191,12 +199,16 @@ public abstract class AbstractMessageSource extends AbstractExpressionEvaluat .build(); } if (this.countsEnabled && message != null) { - if (this.counter != null) { - this.counter.increment(); - } - else { - this.messageCount.incrementAndGet(); + if (this.meterRegistry != null) { + Counter.builder(RECEIVE_COUNTER_NAME) + .tag("name", getComponentName() == null ? "unknown" : getComponentName()) + .tag("type", "source") + .tag("result", "success") + .tag("exception", "none") + .description("Messages received") + .register(this.meterRegistry).increment(); } + this.messageCount.incrementAndGet(); } return message; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageHandler.java index 7f496163e9..2987ec445f 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageHandler.java @@ -36,6 +36,9 @@ import org.springframework.messaging.MessageHandlingException; import org.springframework.messaging.MessagingException; import org.springframework.util.Assert; +import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.core.instrument.Timer; +import io.micrometer.core.instrument.Timer.Sample; import reactor.core.CoreSubscriber; /** @@ -71,6 +74,8 @@ public abstract class AbstractMessageHandler extends IntegrationObjectSupport im private volatile boolean loggingEnabled = true; + private MeterRegistry meterRegistry; + @Override public boolean isLoggingEnabled() { return this.loggingEnabled; @@ -82,6 +87,11 @@ public abstract class AbstractMessageHandler extends IntegrationObjectSupport im this.managementOverrides.loggingConfigured = true; } + @Override + public void registerMeterRegistry(MeterRegistry meterRegistry) { + this.meterRegistry = meterRegistry; + } + @Override public void setOrder(int order) { this.order = order; @@ -131,36 +141,44 @@ public abstract class AbstractMessageHandler extends IntegrationObjectSupport im MetricsContext start = null; boolean countsEnabled = this.countsEnabled; AbstractMessageHandlerMetrics handlerMetrics = this.handlerMetrics; + Sample sample = null; + if (countsEnabled && this.meterRegistry != null) { + sample = Timer.start(this.meterRegistry); + } try { if (this.shouldTrack) { message = MessageHistory.write(message, this, this.getMessageBuilderFactory()); } if (countsEnabled) { - if (handlerMetrics.getTimer() != null) { - final Message messageToSend = message; - handlerMetrics.getTimer().recordCallable(() -> { - handleMessageInternal(messageToSend); - return null; - }); - } - else { - start = handlerMetrics.beforeHandle(); - handleMessageInternal(message); - handlerMetrics.afterHandle(start, true); + start = handlerMetrics.beforeHandle(); + handleMessageInternal(message); + if (this.meterRegistry != null) { + sample.stop(Timer.builder(SEND_TIMER_NAME) + .tag("type", "handler") + .tag("name", getComponentName() == null ? "unknown" : getComponentName()) + .tag("result", "success") + .tag("exception", "none") + .description("Subflow process time") + .register(this.meterRegistry)); } + handlerMetrics.afterHandle(start, true); } else { handleMessageInternal(message); } } catch (Exception e) { + if (sample != null) { + sample.stop(Timer.builder(SEND_TIMER_NAME) + .tag("type", "handler") + .tag("name", getComponentName() == null ? "unknown" : getComponentName()) + .tag("result", "failure") + .tag("exception", e.getClass().getSimpleName()) + .description("Subflow process time") + .register(this.meterRegistry)); + } if (countsEnabled) { - if (handlerMetrics.getErrorCounter() != null) { - handlerMetrics.getErrorCounter().increment(); - } - else { - handlerMetrics.afterHandle(start, false); - } + handlerMetrics.afterHandle(start, false); } if (e instanceof MessagingException) { throw (MessagingException) e; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagement.java b/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagement.java index cd57c11459..dcc3c41264 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagement.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagement.java @@ -19,6 +19,8 @@ package org.springframework.integration.support.management; import org.springframework.jmx.export.annotation.ManagedAttribute; import org.springframework.jmx.export.annotation.ManagedOperation; +import io.micrometer.core.instrument.MeterRegistry; + /** * Base interface for Integration managed components. * @@ -28,6 +30,12 @@ import org.springframework.jmx.export.annotation.ManagedOperation; */ public interface IntegrationManagement { + String METER_PREFIX = "spring.integration."; + + String SEND_TIMER_NAME = METER_PREFIX + "send"; + + String RECEIVE_COUNTER_NAME = METER_PREFIX + "receive"; + @ManagedAttribute(description = "Use to disable debug logging during normal message flow") void setLoggingEnabled(boolean enabled); @@ -50,6 +58,15 @@ public interface IntegrationManagement { */ ManagementOverrides getOverrides(); + /** + * Inject a micrometer {@link MeterRegistry} + * @param registry the registry. + * @since 5.0.3 + */ + default void registerMeterRegistry(MeterRegistry registry) { + // no op + } + /** * Toggles to inform the management configurer to not set these properties since * the user has manually configured them in a bean definition. If true, the diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagementConfigurer.java b/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagementConfigurer.java index e475b79901..edaf3f068e 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagementConfigurer.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/management/IntegrationManagementConfigurer.java @@ -26,15 +26,22 @@ import org.apache.commons.logging.LogFactory; import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanNameAware; +import org.springframework.beans.factory.NoSuchBeanDefinitionException; import org.springframework.beans.factory.SmartInitializingSingleton; import org.springframework.beans.factory.config.BeanPostProcessor; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; +import org.springframework.integration.core.MessageSource; import org.springframework.integration.support.management.IntegrationManagement.ManagementOverrides; import org.springframework.integration.util.PatternMatchUtils; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.MessageHandler; import org.springframework.util.Assert; import org.springframework.util.StringUtils; +import io.micrometer.core.instrument.Gauge; +import io.micrometer.core.instrument.MeterRegistry; + /** * Configures beans that implement {@link IntegrationManagement}. @@ -82,6 +89,8 @@ public class IntegrationManagementConfigurer implements SmartInitializingSinglet private volatile boolean singletonsInstantiated; + private MeterRegistry meterRegistry; + @Override public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { this.applicationContext = applicationContext; @@ -205,6 +214,16 @@ public class IntegrationManagementConfigurer implements SmartInitializingSinglet Assert.state(this.applicationContext != null, "'applicationContext' must not be null"); Assert.state(MANAGEMENT_CONFIGURER_NAME.equals(this.beanName), getClass().getSimpleName() + " bean name must be " + MANAGEMENT_CONFIGURER_NAME); + try { + this.meterRegistry = this.applicationContext.getBean(MeterRegistry.class); + } + catch (NoSuchBeanDefinitionException e) { + // no op + } + if (this.meterRegistry != null) { + injectRegistry(this.meterRegistry); + registerComponentGauges(this.meterRegistry); + } if (this.metricsFactory == null && StringUtils.hasText(this.metricsFactoryBeanName)) { this.metricsFactory = this.applicationContext.getBean(this.metricsFactoryBeanName, MetricsFactory.class); } @@ -230,9 +249,26 @@ public class IntegrationManagementConfigurer implements SmartInitializingSinglet this.singletonsInstantiated = true; } + /** + * @param registry + */ + private void injectRegistry(MeterRegistry registry) { + Map managed = this.applicationContext.getBeansOfType(IntegrationManagement.class); + for (Entry entry : managed.entrySet()) { + IntegrationManagement bean = entry.getValue(); + if (!bean.getOverrides().loggingConfigured) { + bean.setLoggingEnabled(this.defaultLoggingEnabled); + } + bean.registerMeterRegistry(registry); + } + } + @Override public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException { if (this.singletonsInstantiated) { + if (bean instanceof IntegrationManagement) { + ((IntegrationManagement) bean).registerMeterRegistry(this.meterRegistry); + } return doConfigureMetrics(bean, beanName); } return bean; @@ -334,6 +370,23 @@ public class IntegrationManagementConfigurer implements SmartInitializingSinglet this.sourcesByName.put(bean.getManagedName() != null ? bean.getManagedName() : name, bean); } + private void registerComponentGauges(MeterRegistry meterRegistry) { + Gauge.builder("spring.integration.channels", this, + (c) -> this.applicationContext.getBeansOfType(MessageChannel.class).size()) + .description("The number of message channels") + .register(meterRegistry); + + Gauge.builder("spring.integration.handlers", this, + (c) -> this.applicationContext.getBeansOfType(MessageHandler.class).size()) + .description("The number of message handlers") + .register(meterRegistry); + + Gauge.builder("spring.integration.sources", this, + (c) -> this.applicationContext.getBeansOfType(MessageSource.class).size()) + .description("The number of message sources") + .register(meterRegistry); + } + public String[] getChannelNames() { return this.channelsByName.keySet().toArray(new String[this.channelsByName.size()]); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsFactory.java b/spring-integration-core/src/main/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsFactory.java index de03dcb7ee..42c7897988 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsFactory.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsFactory.java @@ -48,8 +48,11 @@ import io.micrometer.core.instrument.MeterRegistry; * * @since 5.0.2 * + * @deprecated - micrometer metrics are now in-built. + * * @see org.springframework.integration.support.management.IntegrationManagementConfigurer */ +@Deprecated public class MicrometerMetricsFactory implements MetricsFactory, MessageSourceMetricsConfigurer, ApplicationContextAware, SmartInitializingSingleton { diff --git a/spring-integration-core/src/test/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsTests.java b/spring-integration-core/src/test/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsTests.java index 95f4dd7c5f..f3b3c35f9f 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/support/management/micrometer/MicrometerMetricsTests.java @@ -19,8 +19,6 @@ package org.springframework.integration.support.management.micrometer; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.fail; -import java.util.List; - import org.junit.Test; import org.junit.runner.RunWith; @@ -39,7 +37,6 @@ import org.springframework.integration.config.EnableIntegration; import org.springframework.integration.config.EnableIntegrationManagement; import org.springframework.integration.core.MessageSource; import org.springframework.integration.endpoint.AbstractMessageSource; -import org.springframework.integration.test.util.TestUtils; import org.springframework.messaging.Message; import org.springframework.messaging.MessagingException; import org.springframework.messaging.PollableChannel; @@ -47,11 +44,7 @@ import org.springframework.messaging.support.GenericMessage; import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.junit4.SpringRunner; -import io.micrometer.core.instrument.Counter; -import io.micrometer.core.instrument.Gauge; -import io.micrometer.core.instrument.Meter; import io.micrometer.core.instrument.MeterRegistry; -import io.micrometer.core.instrument.Timer; import io.micrometer.core.instrument.simple.SimpleMeterRegistry; /** @@ -104,70 +97,57 @@ public class MicrometerMetricsTests { catch (RuntimeException e) { assertThat(e.getMessage()).isEqualTo("badPoll"); } - List meters = this.meterRegistry.getMeters(); - assertThat(meters.size()).isEqualTo(22); - for (Meter meter : meters) { - String name = meter.getId().getName(); - switch (name) { - case "channel.timer": - case "micrometerMetricsTests.Config.service.serviceActivator.handler.timer": - case "queue.timer": - assertThat(((Timer) meter).count()).isEqualTo(2L); - break; - case "badPoll.timer": - assertThat(((Timer) meter).count()).isEqualTo(1L); - break; - case "source.counter": - case "channel.errorCounter": - case "micrometerMetricsTests.Config.service.serviceActivator.handler.errorCounter": - assertThat(((Counter) meter).count()).isEqualTo(1L); - break; - case "queue.receive.counter": - assertThat(((Counter) meter).count()).isEqualTo(1L); - break; - case "spring.integration.channels": - assertThat(((Gauge) meter).measure().iterator().next().getValue()).isEqualTo(5.0); - break; - case "spring.integration.handlers": - assertThat(((Gauge) meter).measure().iterator().next().getValue()).isEqualTo(2.0); - break; - case "spring.integration.sources": - assertThat(((Gauge) meter).measure().iterator().next().getValue()).isEqualTo(1.0); - break; - case "errorChannel.timer": - case "nullChannel.timer": - case "_org.springframework.integration.errorLogger.handler.timer": - assertThat(((Timer) meter).count()).isEqualTo(0L); - break; - case "queue.errorCounter": - case "nullChannel.errorCounter": - case "queue.receive.errorCounter": - case "_org.springframework.integration.errorLogger.handler.errorCounter": - case "errorChannel.errorCounter": - case "badPoll.errorCounter": - case "badPoll.receiveCounter": - assertThat(((Counter) meter).count()).isEqualTo(0L); - break; - case "badPoll.receiveErrorCounter": - assertThat(((Counter) meter).count()).isEqualTo(1L); - break; - default: - } - } + MeterRegistry registry = this.meterRegistry; + assertThat(registry.get("spring.integration.channels").gauge().value()).isEqualTo(5); + assertThat(registry.get("spring.integration.handlers").gauge().value()).isEqualTo(2); + assertThat(registry.get("spring.integration.sources").gauge().value()).isEqualTo(1); + + assertThat(registry.get("spring.integration.receive") + .tag("name", "source") + .tag("result", "success") + .counter().count()).isEqualTo(1); + + assertThat(registry.get("spring.integration.receive") + .tag("name", "badPoll") + .tag("result", "failure") + .counter().count()).isEqualTo(1); + + assertThat(registry.get("spring.integration.send") + .tag("name", "micrometerMetricsTests.Config.service.serviceActivator.handler") + .tag("result", "success") + .timer().count()).isEqualTo(1); + + assertThat(registry.get("spring.integration.send") + .tag("name", "channel") + .tag("result", "success") + .timer().count()).isEqualTo(1); + + assertThat(registry.get("spring.integration.send") + .tag("name", "channel") + .tag("result", "failure") + .timer().count()).isEqualTo(1); + + assertThat(registry.get("spring.integration.send") + .tag("name", "micrometerMetricsTests.Config.service.serviceActivator.handler") + .tag("result", "failure") + .timer().count()).isEqualTo(1); + + assertThat(registry.get("spring.integration.receive") + .tag("name", "queue") + .tag("result", "success") + .counter().count()).isEqualTo(1); + BeanDefinitionRegistry beanFactory = (BeanDefinitionRegistry) this.context.getBeanFactory(); beanFactory.registerBeanDefinition("newChannel", BeanDefinitionBuilder.genericBeanDefinition(DirectChannel.class).getRawBeanDefinition()); DirectChannel newChannel = this.context.getBean("newChannel", DirectChannel.class); - assertThat(this.meterRegistry.getMeters().size()).isEqualTo(24); - Timer timer = meterRegistry.get("newChannel.timer").timer(); - assertThat(timer).isSameAs(TestUtils.getPropertyValue(newChannel, "channelMetrics.timer")); - beanFactory.removeBeanDefinition("newChannel"); - // verify that the meter registry reuses the existing timer - beanFactory.registerBeanDefinition("newChannel", - BeanDefinitionBuilder.genericBeanDefinition(DirectChannel.class).getRawBeanDefinition()); - newChannel = this.context.getBean("newChannel", DirectChannel.class); - assertThat(this.meterRegistry.getMeters().size()).isEqualTo(24); - assertThat(timer).isSameAs(TestUtils.getPropertyValue(newChannel, "channelMetrics.timer")); + newChannel.setBeanName("newChannel"); + newChannel.subscribe(m -> { }); + newChannel.send(new GenericMessage<>("foo")); + assertThat(registry.get("spring.integration.send") + .tag("name", "newChannel") + .tag("result", "success") + .timer().count()).isEqualTo(1); } @Configuration @@ -175,11 +155,6 @@ public class MicrometerMetricsTests { @EnableIntegrationManagement public static class Config { - @Bean - public MicrometerMetricsFactory metricsFactory(MeterRegistry meterRegistry) { - return new MicrometerMetricsFactory(meterRegistry); - } - @Bean public MeterRegistry meterRegistry() { return new SimpleMeterRegistry();