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();