diff --git a/spring-cloud-stream-schema/foodorder.avro b/spring-cloud-stream-schema/foodorder.avro new file mode 100644 index 000000000..ca3f36d4d Binary files /dev/null and b/spring-cloud-stream-schema/foodorder.avro differ diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindersHealthIndicatorAutoConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindersHealthIndicatorAutoConfiguration.java index d1e3ff0fe..ff8a5f857 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindersHealthIndicatorAutoConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindersHealthIndicatorAutoConfiguration.java @@ -16,17 +16,17 @@ package org.springframework.cloud.stream.config; +import java.util.Iterator; +import java.util.LinkedHashMap; import java.util.Map; -import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.boot.actuate.autoconfigure.endpoint.EndpointAutoConfiguration; import org.springframework.boot.actuate.autoconfigure.health.ConditionalOnEnabledHealthIndicator; -import org.springframework.boot.actuate.health.AbstractHealthIndicator; -import org.springframework.boot.actuate.health.CompositeHealthIndicator; -import org.springframework.boot.actuate.health.DefaultHealthIndicatorRegistry; +import org.springframework.boot.actuate.health.CompositeHealthContributor; import org.springframework.boot.actuate.health.Health; +import org.springframework.boot.actuate.health.HealthContributor; import org.springframework.boot.actuate.health.HealthIndicator; -import org.springframework.boot.actuate.health.OrderedHealthAggregator; +import org.springframework.boot.actuate.health.NamedContributor; import org.springframework.boot.autoconfigure.AutoConfigureAfter; import org.springframework.boot.autoconfigure.AutoConfigureBefore; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; @@ -50,58 +50,76 @@ import org.springframework.context.annotation.Configuration; public class BindersHealthIndicatorAutoConfiguration { @Bean - @ConditionalOnMissingBean(name = "bindersHealthIndicator") - public CompositeHealthIndicator bindersHealthIndicator() { - return new CompositeHealthIndicator(new OrderedHealthAggregator(), - new DefaultHealthIndicatorRegistry()); + @ConditionalOnMissingBean + public BindersHealthContributor bindersHealthContributor() { + return new BindersHealthContributor(); } @Bean public DefaultBinderFactory.Listener bindersHealthIndicatorListener( - @Qualifier("bindersHealthIndicator") CompositeHealthIndicator compositeHealthIndicator) { - return new BindersHealthIndicatorListener(compositeHealthIndicator); + BindersHealthContributor bindersHealthContributor) { + return new BindersHealthIndicatorListener(bindersHealthContributor); } /** * A {@link DefaultBinderFactory.Listener} that provides {@link HealthIndicator} * support. - * - * @author Ilayaperumal Gopinathan */ private static class BindersHealthIndicatorListener implements DefaultBinderFactory.Listener { - private final CompositeHealthIndicator bindersHealthIndicator; + private final BindersHealthContributor bindersHealthContributor; - BindersHealthIndicatorListener(CompositeHealthIndicator bindersHealthIndicator) { - this.bindersHealthIndicator = bindersHealthIndicator; + BindersHealthIndicatorListener(BindersHealthContributor bindersHealthContributor) { + this.bindersHealthContributor = bindersHealthContributor; } @Override public void afterBinderContextInitialized(String binderConfigurationName, ConfigurableApplicationContext binderContext) { - if (this.bindersHealthIndicator != null) { - OrderedHealthAggregator healthAggregator = new OrderedHealthAggregator(); - Map indicators = binderContext - .getBeansOfType(HealthIndicator.class); - // if there are no health indicators in the child context, we just mark - // the binder's health as unknown - // this can happen due to the fact that configuration is inherited - HealthIndicator binderHealthIndicator = indicators.isEmpty() - ? new DefaultHealthIndicator() - : new CompositeHealthIndicator(healthAggregator, indicators); - this.bindersHealthIndicator.getRegistry() - .register(binderConfigurationName, binderHealthIndicator); + if (this.bindersHealthContributor != null) { + this.bindersHealthContributor.add(binderConfigurationName, + binderContext.getBeansOfType(HealthContributor.class)); } } - private static class DefaultHealthIndicator extends AbstractHealthIndicator { + } - @Override - protected void doHealthCheck(Health.Builder builder) throws Exception { - builder.unknown(); + /** + * {@link CompositeHealthContributor} that provides binder health contributions. + */ + private static class BindersHealthContributor implements CompositeHealthContributor { + + private static final HealthIndicator UNKNOWN = () -> Health.unknown().build(); + + private Map contributors = new LinkedHashMap<>(); + + void add(String binderConfigurationName, Map binderHealthContributors) { + // if there are no health contributors in the child context, we just mark + // the binder's health as unknown + // this can happen due to the fact that configuration is inherited + this.contributors.put(binderConfigurationName, getContributor(binderHealthContributors)); + } + + private HealthContributor getContributor(Map binderHealthContributors) { + if (binderHealthContributors.isEmpty()) { + return UNKNOWN; } + if (binderHealthContributors.size() == 1) { + return binderHealthContributors.values().iterator().next(); + } + return CompositeHealthContributor.fromMap(binderHealthContributors); + } + @Override + public HealthContributor getContributor(String name) { + return contributors.get(name); + } + + @Override + public Iterator> iterator() { + return contributors.entrySet().stream() + .map((entry) -> NamedContributor.of(entry.getKey(), entry.getValue())).iterator(); } } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/HealthIndicatorsConfigurationTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/HealthIndicatorsConfigurationTests.java index 80999b1d1..6b81ae83a 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/HealthIndicatorsConfigurationTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/HealthIndicatorsConfigurationTests.java @@ -19,18 +19,16 @@ package org.springframework.cloud.stream.binder; import java.io.IOException; import java.net.URL; import java.net.URLClassLoader; -import java.util.Map; import org.junit.Test; -import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.NoSuchBeanDefinitionException; import org.springframework.boot.WebApplicationType; -import org.springframework.boot.actuate.health.CompositeHealthIndicator; -import org.springframework.boot.actuate.health.DefaultHealthIndicatorRegistry; +import org.springframework.boot.actuate.health.CompositeHealthContributor; +import org.springframework.boot.actuate.health.Health; +import org.springframework.boot.actuate.health.HealthContributor; import org.springframework.boot.actuate.health.HealthIndicator; -import org.springframework.boot.actuate.health.HealthIndicatorRegistry; -import org.springframework.boot.actuate.health.OrderedHealthAggregator; +import org.springframework.boot.actuate.health.NamedContributor; import org.springframework.boot.actuate.health.Status; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.builder.SpringApplicationBuilder; @@ -86,28 +84,22 @@ public class HealthIndicatorsConfigurationTests { Binder binder2 = context.getBean(BinderFactory.class).getBinder("binder2", MessageChannel.class); assertThat(binder2).isInstanceOf(StubBinder2.class); - CompositeHealthIndicator bindersHealthIndicator = context - .getBean("bindersHealthIndicator", CompositeHealthIndicator.class); - DirectFieldAccessor directFieldAccessor = new DirectFieldAccessor( - bindersHealthIndicator); - assertThat(bindersHealthIndicator).isNotNull(); + CompositeHealthContributor bindersHealthContributor = context + .getBean("bindersHealthContributor", CompositeHealthContributor.class); + assertThat(bindersHealthContributor).isNotNull(); assertThat( - context.getBean("test1HealthIndicator1", CompositeHealthIndicator.class)) + context.getBean("test1HealthIndicator1", HealthContributor.class)) .isNotNull(); assertThat( - context.getBean("test2HealthIndicator2", CompositeHealthIndicator.class)) + context.getBean("test2HealthIndicator2", HealthContributor.class)) .isNotNull(); - HealthIndicatorRegistry registry = (HealthIndicatorRegistry) directFieldAccessor - .getPropertyValue("registry"); - - Map healthIndicators = registry.getAll(); - assertThat(healthIndicators).containsKey("binder1"); - assertThat(healthIndicators.get("binder1").health().getStatus()) + assertThat(bindersHealthContributor.stream().map(NamedContributor::getName)).contains("binder1", "binder2"); + assertThat(bindersHealthContributor.getContributor("binder1")).extracting("health").extracting("status") .isEqualTo(Status.UP); - assertThat(healthIndicators).containsKey("binder2"); - assertThat(healthIndicators.get("binder2").health().getStatus()) + assertThat(bindersHealthContributor.getContributor("binder2")).extracting("health").extracting("status") .isEqualTo(Status.UNKNOWN); + context.close(); } @@ -126,16 +118,16 @@ public class HealthIndicatorsConfigurationTests { MessageChannel.class); assertThat(binder2).isInstanceOf(StubBinder2.class); try { - context.getBean("bindersHealthIndicator", CompositeHealthIndicator.class); - fail("The 'bindersHealthIndicator' bean should have not been defined"); + context.getBean("bindersHealthContributor", CompositeHealthContributor.class); + fail("The 'bindersHealthContributor' bean should have not been defined"); } catch (NoSuchBeanDefinitionException e) { } assertThat( - context.getBean("test1HealthIndicator1", CompositeHealthIndicator.class)) + context.getBean("test1HealthIndicator1", HealthContributor.class)) .isNotNull(); assertThat( - context.getBean("test2HealthIndicator2", CompositeHealthIndicator.class)) + context.getBean("test2HealthIndicator2", HealthContributor.class)) .isNotNull(); context.close(); } @@ -148,13 +140,13 @@ public class HealthIndicatorsConfigurationTests { static class TestConfig { @Bean - public CompositeHealthIndicator test1HealthIndicator1() { - return new CompositeHealthIndicator(new OrderedHealthAggregator(), new DefaultHealthIndicatorRegistry()); + public HealthIndicator test1HealthIndicator1() { + return () -> Health.unknown().build(); } @Bean - public CompositeHealthIndicator test2HealthIndicator2() { - return new CompositeHealthIndicator(new OrderedHealthAggregator(), new DefaultHealthIndicatorRegistry()); + public HealthIndicator test2HealthIndicator2() { + return () -> Health.unknown().build(); } }