diff --git a/spring-cloud-netflix-hystrix-amqp/pom.xml b/spring-cloud-netflix-hystrix-amqp/pom.xml index 4530c8c4..4dc95cb9 100644 --- a/spring-cloud-netflix-hystrix-amqp/pom.xml +++ b/spring-cloud-netflix-hystrix-amqp/pom.xml @@ -76,6 +76,16 @@ compile true + + org.springframework.boot + spring-boot-starter-web + test + + + org.springframework.boot + spring-boot-starter-actuator + test + diff --git a/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamAmqpProperties.java b/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamAmqpProperties.java new file mode 100644 index 00000000..9c06b5c5 --- /dev/null +++ b/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamAmqpProperties.java @@ -0,0 +1,15 @@ +package org.springframework.netflix.hystrix.amqp; + +import lombok.Data; +import org.springframework.boot.context.properties.ConfigurationProperties; + +/** + * @author Spencer Gibb + */ +@ConfigurationProperties("hystrix.stream.amqp") +@Data +public class HystrixStreamAmqpProperties { + private boolean enabled = true; + private boolean prefixMetricName = true; + private boolean sendId = true; +} diff --git a/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamAutoConfiguration.java b/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamAutoConfiguration.java index 9d2cb054..056cf349 100644 --- a/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamAutoConfiguration.java +++ b/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamAutoConfiguration.java @@ -8,6 +8,7 @@ import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; +import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.netflix.Constants; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -28,6 +29,7 @@ import com.netflix.hystrix.HystrixCircuitBreaker; @ConditionalOnClass({ HystrixCircuitBreaker.class, RabbitTemplate.class }) @ConditionalOnExpression("${hystrix.stream.amqp.enabled:true}") @IntegrationComponentScan(basePackageClasses = HystrixStreamChannel.class) +@EnableConfigurationProperties @EnableScheduling public class HystrixStreamAutoConfiguration { @@ -37,12 +39,18 @@ public class HystrixStreamAutoConfiguration { @Autowired(required = false) private ObjectMapper objectMapper; + @PostConstruct public void init() { Jackson2JsonMessageConverter converter = messageConverter(); amqpTemplate.setMessageConverter(converter); } + @Bean + public HystrixStreamAmqpProperties hystrixStreamAmqpProperties() { + return new HystrixStreamAmqpProperties(); + } + @Bean public HystrixStreamTask hystrixStreamTask() { return new HystrixStreamTask(); diff --git a/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamTask.java b/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamTask.java index bb42f6fe..83d42684 100644 --- a/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamTask.java +++ b/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamTask.java @@ -5,10 +5,12 @@ import com.netflix.hystrix.util.HystrixRollingNumberEvent; import lombok.extern.slf4j.Slf4j; import org.codehaus.jackson.JsonFactory; import org.codehaus.jackson.JsonGenerator; +import org.springframework.beans.BeansException; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cloud.client.ServiceInstance; import org.springframework.cloud.client.discovery.DiscoveryClient; import org.springframework.context.ApplicationContext; +import org.springframework.context.ApplicationContextAware; import org.springframework.scheduling.annotation.Scheduled; import java.io.IOException; @@ -22,7 +24,7 @@ import java.util.concurrent.LinkedBlockingQueue; * see com.netflix.hystrix.contrib.metrics.eventstream.HystrixMetricsPoller.MetricsPoller */ @Slf4j -public class HystrixStreamTask { +public class HystrixStreamTask implements ApplicationContextAware { @Autowired private HystrixStreamChannel channel; @@ -30,13 +32,20 @@ public class HystrixStreamTask { @Autowired private DiscoveryClient discoveryClient; - @Autowired private ApplicationContext context; + @Autowired + private HystrixStreamAmqpProperties properties; + private final LinkedBlockingQueue jsonMetrics = new LinkedBlockingQueue<>(1000); private final JsonFactory jsonFactory = new JsonFactory(); + @Override + public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { + this.context = applicationContext; + } + //TODO: use integration to split this up? @Scheduled(fixedRateString = "${hystrix.stream.amqp.sendRate:500}") public void sendMetrics() { @@ -66,6 +75,9 @@ public class HystrixStreamTask { if (!instances.isEmpty()) { log.trace("gathering metrics size: " + instances.size()); } + + ServiceInstance localService = discoveryClient.getLocalServiceInstance(); + for (HystrixCommandMetrics commandMetrics : instances) { HystrixCommandKey key = commandMetrics.getCommandKey(); HystrixCircuitBreaker circuitBreaker = HystrixCircuitBreaker.Factory.getInstance(key); @@ -73,12 +85,19 @@ public class HystrixStreamTask { StringWriter jsonString = new StringWriter(); JsonGenerator json = jsonFactory.createJsonGenerator(jsonString); + json.writeStartObject(); - addServiceData(json); + addServiceData(json, localService); json.writeObjectFieldStart("data"); json.writeStringField("type", "HystrixCommand"); - json.writeStringField("name", key.name()); + String name = key.name(); + + if (properties.isPrefixMetricName()) { + name = localService.getServiceId() +"."+name; + } + + json.writeStringField("name", name); json.writeStringField("group", commandMetrics.getCommandGroup().name()); json.writeNumberField("currentTime", System.currentTimeMillis()); @@ -187,7 +206,7 @@ public class HystrixStreamTask { JsonGenerator json = jsonFactory.createJsonGenerator(jsonString); json.writeStartObject(); - addServiceData(json); + addServiceData(json, localService); json.writeObjectFieldStart("data"); json.writeStringField("type", "HystrixThreadPool"); @@ -221,13 +240,15 @@ public class HystrixStreamTask { } } - private void addServiceData(JsonGenerator json) throws IOException { - ServiceInstance localService = discoveryClient.getLocalServiceInstance(); + private void addServiceData(JsonGenerator json, ServiceInstance localService) throws IOException { + json.writeObjectFieldStart("origin"); json.writeStringField("host", localService.getHost()); json.writeNumberField("port", localService.getPort()); json.writeStringField("serviceId", localService.getServiceId()); - json.writeStringField("id", context.getId()); + if (properties.isSendId()) { + json.writeStringField("id", context.getId()); + } json.writeEndObject(); } }