diff --git a/.travis.yml b/.travis.yml
index 1109f4c5..bf40db2e 100644
--- a/.travis.yml
+++ b/.travis.yml
@@ -3,6 +3,8 @@ cache:
directories:
- $HOME/.m2
language: java
+jdk:
+ - oraclejdk8
before_install:
- git config user.name "$GIT_NAME"
- git config user.email "$GIT_EMAIL"
diff --git a/pom.xml b/pom.xml
index f6cdb187..95c124cd 100644
--- a/pom.xml
+++ b/pom.xml
@@ -28,8 +28,10 @@
spring-cloud-netflix-core
spring-cloud-netflix-hystrix-dashboard
+ spring-cloud-netflix-hystrix-amqp
spring-cloud-netflix-eureka-server
spring-cloud-netflix-turbine
+ spring-cloud-netflix-turbine-amqp
spring-cloud-netflix-sidecar
docs
@@ -184,21 +186,6 @@
rxjava-core
${netflix.rxjava.version}
-
- com.netflix.turbine
- turbine-core
- ${turbine.version}
-
-
- javax.servlet
- servlet-api
-
-
- log4j
- log4j
-
-
-
com.netflix.zuul
zuul-core
@@ -210,6 +197,11 @@
+
+ org.springframework.integration
+ spring-integration-java-dsl
+ ${spring-integration-dsl.version}
+
org.projectlombok
lombok
@@ -265,10 +257,11 @@
6.1.3
1.4.0-RC5
2.0-RC13
- 1.0.0
1.0.28
0.20.7
1.7
+ 1.0.0.RELEASE
+ 1.1.1.BUILD-SNAPSHOT
diff --git a/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/Constants.java b/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/Constants.java
new file mode 100644
index 00000000..77d7bef4
--- /dev/null
+++ b/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/Constants.java
@@ -0,0 +1,8 @@
+package org.springframework.cloud.netflix;
+
+/**
+ * @author Spencer Gibb
+ */
+public interface Constants {
+ String HYSTRIX_STREAM_NAME = "spring.cloud.hystrix.stream";
+}
diff --git a/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/hystrix/HystrixCircuitBreakerConfiguration.java b/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/hystrix/HystrixCircuitBreakerConfiguration.java
index b2e0ae6d..d2441d13 100644
--- a/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/hystrix/HystrixCircuitBreakerConfiguration.java
+++ b/spring-cloud-netflix-core/src/main/java/org/springframework/cloud/netflix/hystrix/HystrixCircuitBreakerConfiguration.java
@@ -21,11 +21,13 @@ import com.netflix.hystrix.Hystrix;
import com.netflix.hystrix.contrib.javanica.aop.aspectj.HystrixCommandAspect;
import com.netflix.hystrix.contrib.metrics.eventstream.HystrixMetricsPoller;
import com.netflix.hystrix.contrib.metrics.eventstream.HystrixMetricsPoller.MetricsAsJsonPollerListener;
+import com.netflix.hystrix.contrib.metrics.eventstream.HystrixMetricsStreamServlet;
import org.apache.catalina.core.ApplicationContext;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.actuate.endpoint.Endpoint;
import org.springframework.boot.actuate.metrics.GaugeService;
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
@@ -61,6 +63,7 @@ public class HystrixCircuitBreakerConfiguration {
@Configuration
@ConditionalOnExpression("${hystrix.stream.endpoint.enabled:true}")
@ConditionalOnWebApplication
+ @ConditionalOnClass({Endpoint.class, HystrixMetricsStreamServlet.class})
protected static class HystrixWebConfiguration {
@Bean
@@ -71,7 +74,7 @@ public class HystrixCircuitBreakerConfiguration {
}
@Configuration
- @ConditionalOnClass(GaugeService.class)
+ @ConditionalOnClass({HystrixMetricsPoller.class, GaugeService.class})
protected static class HystrixMetricsPollerConfiguration implements SmartLifecycle {
private static Log logger = LogFactory
diff --git a/spring-cloud-netflix-hystrix-amqp/pom.xml b/spring-cloud-netflix-hystrix-amqp/pom.xml
new file mode 100644
index 00000000..531c7911
--- /dev/null
+++ b/spring-cloud-netflix-hystrix-amqp/pom.xml
@@ -0,0 +1,77 @@
+
+
+ 4.0.0
+
+ spring-cloud-netflix-hystrix-amqp
+ jar
+
+ Spring Cloud Netflix Hystrix AMQP
+ Spring Cloud Netflix Hystrix AMQP
+
+
+ org.springframework.cloud
+ spring-cloud-netflix
+ 1.0.0.BUILD-SNAPSHOT
+ ..
+
+
+
+
+ org.springframework.cloud
+ spring-cloud-commons
+
+
+ org.springframework.cloud
+ spring-cloud-netflix-core
+
+
+ org.springframework.boot
+ spring-boot-starter-amqp
+
+
+ org.springframework.boot
+ spring-boot-starter-integration
+
+
+ org.springframework.cloud
+ spring-cloud-netflix-core
+
+
+ org.springframework.integration
+ spring-integration-amqp
+ true
+
+
+ org.springframework.integration
+ spring-integration-java-dsl
+ true
+
+
+ com.netflix.hystrix
+ hystrix-core
+
+
+ com.netflix.hystrix
+ hystrix-metrics-event-stream
+ test
+
+
+ com.netflix.hystrix
+ hystrix-javanica
+ test
+
+
+ com.netflix.eureka
+ eureka-client
+ test
+
+
+ org.projectlombok
+ lombok
+ compile
+ 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
new file mode 100644
index 00000000..bdd0f340
--- /dev/null
+++ b/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamAutoConfiguration.java
@@ -0,0 +1,65 @@
+package org.springframework.netflix.hystrix.amqp;
+
+import com.netflix.hystrix.HystrixCircuitBreaker;
+import org.springframework.amqp.core.AmqpTemplate;
+import org.springframework.amqp.core.DirectExchange;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
+import org.springframework.cloud.netflix.Constants;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.integration.annotation.IntegrationComponentScan;
+import org.springframework.integration.channel.DirectChannel;
+import org.springframework.integration.dsl.IntegrationFlow;
+import org.springframework.integration.dsl.IntegrationFlows;
+import org.springframework.integration.dsl.amqp.Amqp;
+import org.springframework.scheduling.annotation.EnableScheduling;
+
+/**
+ * @author Spencer Gibb
+ */
+@Configuration
+@ConditionalOnClass({HystrixCircuitBreaker.class, AmqpTemplate.class})
+public class HystrixStreamAutoConfiguration {
+ @Configuration
+ @ConditionalOnExpression("${hystrix.stream.amqp.enabled:true}")
+ @IntegrationComponentScan(basePackageClasses = HystrixStreamChannel.class)
+ @EnableScheduling
+ protected static class HystrixStreamBusAutoConfiguration {
+
+ @Autowired
+ private AmqpTemplate amqpTemplate;
+
+ @Bean
+ public HystrixStreamTask hystrixStreamTask() {
+ return new HystrixStreamTask();
+ }
+
+ @Bean
+ public DirectChannel hystrixStream() {
+ return new DirectChannel();
+ }
+
+ @Bean
+ public DirectExchange hystrixStreamExchange() {
+ DirectExchange exchange = new DirectExchange(Constants.HYSTRIX_STREAM_NAME);
+ return exchange;
+ }
+
+ @Bean
+ public IntegrationFlow hystrixStreamOutboundFlow() {
+ return IntegrationFlows.from("hystrixStream")
+ //TODO: set content type
+ /*.enrichHeaders(new ComponentConfigurer() {
+ @Override
+ public void configure(HeaderEnricherSpec spec) {
+ spec.header("content-type", "application/json", true);
+ }
+ })*/
+ .handle(Amqp.outboundAdapter(this.amqpTemplate).exchangeName(Constants.HYSTRIX_STREAM_NAME))
+ .get();
+ }
+ }
+
+}
diff --git a/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamChannel.java b/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamChannel.java
new file mode 100644
index 00000000..95214a61
--- /dev/null
+++ b/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamChannel.java
@@ -0,0 +1,14 @@
+package org.springframework.netflix.hystrix.amqp;
+
+import org.springframework.integration.annotation.Gateway;
+import org.springframework.integration.annotation.MessagingGateway;
+
+/**
+ * @author Spencer Gibb
+ */
+@MessagingGateway
+public interface HystrixStreamChannel {
+
+ @Gateway(requestChannel = "hystrixStream")
+ public void send(String s);
+}
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
new file mode 100644
index 00000000..bb42f6fe
--- /dev/null
+++ b/spring-cloud-netflix-hystrix-amqp/src/main/java/org/springframework/netflix/hystrix/amqp/HystrixStreamTask.java
@@ -0,0 +1,233 @@
+package org.springframework.netflix.hystrix.amqp;
+
+import com.netflix.hystrix.*;
+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.factory.annotation.Autowired;
+import org.springframework.cloud.client.ServiceInstance;
+import org.springframework.cloud.client.discovery.DiscoveryClient;
+import org.springframework.context.ApplicationContext;
+import org.springframework.scheduling.annotation.Scheduled;
+
+import java.io.IOException;
+import java.io.StringWriter;
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.concurrent.LinkedBlockingQueue;
+
+/**
+ * @author Spencer Gibb
+ * see com.netflix.hystrix.contrib.metrics.eventstream.HystrixMetricsPoller.MetricsPoller
+ */
+@Slf4j
+public class HystrixStreamTask {
+
+ @Autowired
+ private HystrixStreamChannel channel;
+
+ @Autowired
+ private DiscoveryClient discoveryClient;
+
+ @Autowired
+ private ApplicationContext context;
+
+ private final LinkedBlockingQueue jsonMetrics = new LinkedBlockingQueue<>(1000);
+
+ private final JsonFactory jsonFactory = new JsonFactory();
+
+ //TODO: use integration to split this up?
+ @Scheduled(fixedRateString = "${hystrix.stream.amqp.sendRate:500}")
+ public void sendMetrics() {
+ ArrayList metrics = new ArrayList<>();
+ jsonMetrics.drainTo(metrics);
+
+ if (!metrics.isEmpty()) {
+ log.trace("sending amqp metrics size: " + metrics.size());
+ for (String json: metrics) {
+ //TODO: batch all metrics to one message
+ try {
+ channel.send(json);
+ } catch (Exception e) {
+ e.printStackTrace();
+ }
+ }
+ }
+ }
+
+ //@InboundChannelAdapter()
+ //TODO: move fixedRate to configuration
+ @Scheduled(fixedRateString = "${hystrix.stream.amqp.gatherRate:500}")
+ public void gatherMetrics() {
+ try {
+ // command metrics
+ Collection instances = HystrixCommandMetrics.getInstances();
+ if (!instances.isEmpty()) {
+ log.trace("gathering metrics size: " + instances.size());
+ }
+ for (HystrixCommandMetrics commandMetrics : instances) {
+ HystrixCommandKey key = commandMetrics.getCommandKey();
+ HystrixCircuitBreaker circuitBreaker = HystrixCircuitBreaker.Factory.getInstance(key);
+
+ StringWriter jsonString = new StringWriter();
+ JsonGenerator json = jsonFactory.createJsonGenerator(jsonString);
+
+ json.writeStartObject();
+
+ addServiceData(json);
+ json.writeObjectFieldStart("data");
+ json.writeStringField("type", "HystrixCommand");
+ json.writeStringField("name", key.name());
+ json.writeStringField("group", commandMetrics.getCommandGroup().name());
+ json.writeNumberField("currentTime", System.currentTimeMillis());
+
+ // circuit breaker
+ if (circuitBreaker == null) {
+ // circuit breaker is disabled and thus never open
+ json.writeBooleanField("isCircuitBreakerOpen", false);
+ } else {
+ json.writeBooleanField("isCircuitBreakerOpen", circuitBreaker.isOpen());
+ }
+ HystrixCommandMetrics.HealthCounts healthCounts = commandMetrics.getHealthCounts();
+ json.writeNumberField("errorPercentage", healthCounts.getErrorPercentage());
+ json.writeNumberField("errorCount", healthCounts.getErrorCount());
+ json.writeNumberField("requestCount", healthCounts.getTotalRequests());
+
+ // rolling counters
+ json.writeNumberField("rollingCountCollapsedRequests", commandMetrics.getRollingCount(HystrixRollingNumberEvent.COLLAPSED));
+ json.writeNumberField("rollingCountExceptionsThrown", commandMetrics.getRollingCount(HystrixRollingNumberEvent.EXCEPTION_THROWN));
+ json.writeNumberField("rollingCountFailure", commandMetrics.getRollingCount(HystrixRollingNumberEvent.FAILURE));
+ json.writeNumberField("rollingCountFallbackFailure", commandMetrics.getRollingCount(HystrixRollingNumberEvent.FALLBACK_FAILURE));
+ json.writeNumberField("rollingCountFallbackRejection", commandMetrics.getRollingCount(HystrixRollingNumberEvent.FALLBACK_REJECTION));
+ json.writeNumberField("rollingCountFallbackSuccess", commandMetrics.getRollingCount(HystrixRollingNumberEvent.FALLBACK_SUCCESS));
+ json.writeNumberField("rollingCountResponsesFromCache", commandMetrics.getRollingCount(HystrixRollingNumberEvent.RESPONSE_FROM_CACHE));
+ json.writeNumberField("rollingCountSemaphoreRejected", commandMetrics.getRollingCount(HystrixRollingNumberEvent.SEMAPHORE_REJECTED));
+ json.writeNumberField("rollingCountShortCircuited", commandMetrics.getRollingCount(HystrixRollingNumberEvent.SHORT_CIRCUITED));
+ json.writeNumberField("rollingCountSuccess", commandMetrics.getRollingCount(HystrixRollingNumberEvent.SUCCESS));
+ json.writeNumberField("rollingCountThreadPoolRejected", commandMetrics.getRollingCount(HystrixRollingNumberEvent.THREAD_POOL_REJECTED));
+ json.writeNumberField("rollingCountTimeout", commandMetrics.getRollingCount(HystrixRollingNumberEvent.TIMEOUT));
+
+ json.writeNumberField("currentConcurrentExecutionCount", commandMetrics.getCurrentConcurrentExecutionCount());
+
+ // latency percentiles
+ json.writeNumberField("latencyExecute_mean", commandMetrics.getExecutionTimeMean());
+ json.writeObjectFieldStart("latencyExecute");
+ json.writeNumberField("0", commandMetrics.getExecutionTimePercentile(0));
+ json.writeNumberField("25", commandMetrics.getExecutionTimePercentile(25));
+ json.writeNumberField("50", commandMetrics.getExecutionTimePercentile(50));
+ json.writeNumberField("75", commandMetrics.getExecutionTimePercentile(75));
+ json.writeNumberField("90", commandMetrics.getExecutionTimePercentile(90));
+ json.writeNumberField("95", commandMetrics.getExecutionTimePercentile(95));
+ json.writeNumberField("99", commandMetrics.getExecutionTimePercentile(99));
+ json.writeNumberField("99.5", commandMetrics.getExecutionTimePercentile(99.5));
+ json.writeNumberField("100", commandMetrics.getExecutionTimePercentile(100));
+ json.writeEndObject();
+ //
+ json.writeNumberField("latencyTotal_mean", commandMetrics.getTotalTimeMean());
+ json.writeObjectFieldStart("latencyTotal");
+ json.writeNumberField("0", commandMetrics.getTotalTimePercentile(0));
+ json.writeNumberField("25", commandMetrics.getTotalTimePercentile(25));
+ json.writeNumberField("50", commandMetrics.getTotalTimePercentile(50));
+ json.writeNumberField("75", commandMetrics.getTotalTimePercentile(75));
+ json.writeNumberField("90", commandMetrics.getTotalTimePercentile(90));
+ json.writeNumberField("95", commandMetrics.getTotalTimePercentile(95));
+ json.writeNumberField("99", commandMetrics.getTotalTimePercentile(99));
+ json.writeNumberField("99.5", commandMetrics.getTotalTimePercentile(99.5));
+ json.writeNumberField("100", commandMetrics.getTotalTimePercentile(100));
+ json.writeEndObject();
+
+ // property values for reporting what is actually seen by the command rather than what was set somewhere
+ HystrixCommandProperties commandProperties = commandMetrics.getProperties();
+
+ json.writeNumberField("propertyValue_circuitBreakerRequestVolumeThreshold", commandProperties.circuitBreakerRequestVolumeThreshold().get());
+ json.writeNumberField("propertyValue_circuitBreakerSleepWindowInMilliseconds", commandProperties.circuitBreakerSleepWindowInMilliseconds().get());
+ json.writeNumberField("propertyValue_circuitBreakerErrorThresholdPercentage", commandProperties.circuitBreakerErrorThresholdPercentage().get());
+ json.writeBooleanField("propertyValue_circuitBreakerForceOpen", commandProperties.circuitBreakerForceOpen().get());
+ json.writeBooleanField("propertyValue_circuitBreakerForceClosed", commandProperties.circuitBreakerForceClosed().get());
+ json.writeBooleanField("propertyValue_circuitBreakerEnabled", commandProperties.circuitBreakerEnabled().get());
+
+ json.writeStringField("propertyValue_executionIsolationStrategy", commandProperties.executionIsolationStrategy().get().name());
+ json.writeNumberField("propertyValue_executionIsolationThreadTimeoutInMilliseconds", commandProperties.executionIsolationThreadTimeoutInMilliseconds().get());
+ json.writeBooleanField("propertyValue_executionIsolationThreadInterruptOnTimeout", commandProperties.executionIsolationThreadInterruptOnTimeout().get());
+ json.writeStringField("propertyValue_executionIsolationThreadPoolKeyOverride", commandProperties.executionIsolationThreadPoolKeyOverride().get());
+ json.writeNumberField("propertyValue_executionIsolationSemaphoreMaxConcurrentRequests", commandProperties.executionIsolationSemaphoreMaxConcurrentRequests().get());
+ json.writeNumberField("propertyValue_fallbackIsolationSemaphoreMaxConcurrentRequests", commandProperties.fallbackIsolationSemaphoreMaxConcurrentRequests().get());
+
+ /*
+ * The following are commented out as these rarely change and are verbose for streaming for something people don't change.
+ * We could perhaps allow a property or request argument to include these.
+ */
+
+ // json.put("propertyValue_metricsRollingPercentileEnabled", commandProperties.metricsRollingPercentileEnabled().get());
+ // json.put("propertyValue_metricsRollingPercentileBucketSize", commandProperties.metricsRollingPercentileBucketSize().get());
+ // json.put("propertyValue_metricsRollingPercentileWindow", commandProperties.metricsRollingPercentileWindowInMilliseconds().get());
+ // json.put("propertyValue_metricsRollingPercentileWindowBuckets", commandProperties.metricsRollingPercentileWindowBuckets().get());
+ // json.put("propertyValue_metricsRollingStatisticalWindowBuckets", commandProperties.metricsRollingStatisticalWindowBuckets().get());
+ json.writeNumberField("propertyValue_metricsRollingStatisticalWindowInMilliseconds", commandProperties.metricsRollingStatisticalWindowInMilliseconds().get());
+
+ json.writeBooleanField("propertyValue_requestCacheEnabled", commandProperties.requestCacheEnabled().get());
+ json.writeBooleanField("propertyValue_requestLogEnabled", commandProperties.requestLogEnabled().get());
+
+ json.writeNumberField("reportingHosts", 1); // this will get summed across all instances in a cluster
+
+ json.writeEndObject(); // end data attribute
+ json.writeEndObject();
+ json.close();
+
+ // output
+ jsonMetrics.add(jsonString.getBuffer().toString());
+ }
+
+ // thread pool metrics
+ for (HystrixThreadPoolMetrics threadPoolMetrics : HystrixThreadPoolMetrics.getInstances()) {
+ HystrixThreadPoolKey key = threadPoolMetrics.getThreadPoolKey();
+
+ StringWriter jsonString = new StringWriter();
+ JsonGenerator json = jsonFactory.createJsonGenerator(jsonString);
+ json.writeStartObject();
+
+ addServiceData(json);
+ json.writeObjectFieldStart("data");
+
+ json.writeStringField("type", "HystrixThreadPool");
+ json.writeStringField("name", key.name());
+ json.writeNumberField("currentTime", System.currentTimeMillis());
+
+ json.writeNumberField("currentActiveCount", threadPoolMetrics.getCurrentActiveCount().intValue());
+ json.writeNumberField("currentCompletedTaskCount", threadPoolMetrics.getCurrentCompletedTaskCount().longValue());
+ json.writeNumberField("currentCorePoolSize", threadPoolMetrics.getCurrentCorePoolSize().intValue());
+ json.writeNumberField("currentLargestPoolSize", threadPoolMetrics.getCurrentLargestPoolSize().intValue());
+ json.writeNumberField("currentMaximumPoolSize", threadPoolMetrics.getCurrentMaximumPoolSize().intValue());
+ json.writeNumberField("currentPoolSize", threadPoolMetrics.getCurrentPoolSize().intValue());
+ json.writeNumberField("currentQueueSize", threadPoolMetrics.getCurrentQueueSize().intValue());
+ json.writeNumberField("currentTaskCount", threadPoolMetrics.getCurrentTaskCount().longValue());
+ json.writeNumberField("rollingCountThreadsExecuted", threadPoolMetrics.getRollingCountThreadsExecuted());
+ json.writeNumberField("rollingMaxActiveThreads", threadPoolMetrics.getRollingMaxActiveThreads());
+
+ json.writeNumberField("propertyValue_queueSizeRejectionThreshold", threadPoolMetrics.getProperties().queueSizeRejectionThreshold().get());
+ json.writeNumberField("propertyValue_metricsRollingStatisticalWindowInMilliseconds", threadPoolMetrics.getProperties().metricsRollingStatisticalWindowInMilliseconds().get());
+
+ json.writeNumberField("reportingHosts", 1); // this will get summed across all instances in a cluster
+
+ json.writeEndObject(); // end of data object
+ json.writeEndObject();
+ json.close();
+ // output to stream
+ jsonMetrics.add(jsonString.getBuffer().toString());
+ }
+ } catch (Exception e) {
+ log.error("Error adding metrics to queue", e);
+ }
+ }
+
+ private void addServiceData(JsonGenerator json) throws IOException {
+ ServiceInstance localService = discoveryClient.getLocalServiceInstance();
+ json.writeObjectFieldStart("origin");
+ json.writeStringField("host", localService.getHost());
+ json.writeNumberField("port", localService.getPort());
+ json.writeStringField("serviceId", localService.getServiceId());
+ json.writeStringField("id", context.getId());
+ json.writeEndObject();
+ }
+}
diff --git a/spring-cloud-netflix-hystrix-amqp/src/main/resources/META-INF/spring.factories b/spring-cloud-netflix-hystrix-amqp/src/main/resources/META-INF/spring.factories
new file mode 100644
index 00000000..9bf3f8f7
--- /dev/null
+++ b/spring-cloud-netflix-hystrix-amqp/src/main/resources/META-INF/spring.factories
@@ -0,0 +1,2 @@
+org.springframework.boot.autoconfigure.EnableAutoConfiguration=\
+org.springframework.netflix.hystrix.amqp.HystrixStreamAutoConfiguration
diff --git a/spring-cloud-netflix-hystrix-amqp/src/test/java/org/springframework/netflix/hystrix/amqp/SampleHystrixAmqpApplicaiton.java b/spring-cloud-netflix-hystrix-amqp/src/test/java/org/springframework/netflix/hystrix/amqp/SampleHystrixAmqpApplicaiton.java
new file mode 100644
index 00000000..9c885aaa
--- /dev/null
+++ b/spring-cloud-netflix-hystrix-amqp/src/test/java/org/springframework/netflix/hystrix/amqp/SampleHystrixAmqpApplicaiton.java
@@ -0,0 +1,29 @@
+package org.springframework.netflix.hystrix.amqp;
+
+import com.netflix.hystrix.contrib.javanica.annotation.HystrixCommand;
+import org.springframework.boot.SpringApplication;
+import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
+import org.springframework.cloud.client.circuitbreaker.EnableCircuitBreaker;
+import org.springframework.cloud.client.discovery.EnableDiscoveryClient;
+import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RestController;
+
+/**
+ * @author Spencer Gibb
+ */
+@EnableAutoConfiguration
+@EnableDiscoveryClient
+@EnableCircuitBreaker
+@RestController
+public class SampleHystrixAmqpApplicaiton {
+
+ @HystrixCommand
+ @RequestMapping("/")
+ public String hello() {
+ return "Hello World";
+ }
+
+ public static void main(String[] args) {
+ SpringApplication.run(SampleHystrixAmqpApplicaiton.class, args);
+ }
+}
diff --git a/spring-cloud-netflix-hystrix-amqp/src/test/resources/application.yml b/spring-cloud-netflix-hystrix-amqp/src/test/resources/application.yml
new file mode 100644
index 00000000..c61d2e27
--- /dev/null
+++ b/spring-cloud-netflix-hystrix-amqp/src/test/resources/application.yml
@@ -0,0 +1,6 @@
+server:
+ port: 17642
+
+logging:
+ level:
+ org.springframework.netflix.hystrix.amqp: TRACE
\ No newline at end of file
diff --git a/spring-cloud-netflix-turbine-amqp/.jdk8 b/spring-cloud-netflix-turbine-amqp/.jdk8
new file mode 100644
index 00000000..e69de29b
diff --git a/spring-cloud-netflix-turbine-amqp/pom.xml b/spring-cloud-netflix-turbine-amqp/pom.xml
new file mode 100644
index 00000000..cdcdc746
--- /dev/null
+++ b/spring-cloud-netflix-turbine-amqp/pom.xml
@@ -0,0 +1,108 @@
+
+
+ 4.0.0
+
+ spring-cloud-netflix-turbine-amqp
+ jar
+
+ Spring Cloud Netflix Turbine AMQP
+ Spring Cloud Netflix Turbine AMQP
+
+
+ org.springframework.cloud
+ spring-cloud-netflix
+ 1.0.0.BUILD-SNAPSHOT
+ ..
+
+
+
+
+
+ org.apache.maven.plugins
+ maven-compiler-plugin
+
+ 1.8
+ 1.8
+
+
+
+
+
+
+ [1.0,1.1)
+ 2.0.0-DP.2
+
+
+
+
+
+ io.reactivex
+ rxjava
+ ${rxjava.version}
+
+
+ com.netflix.turbine
+ turbine-core
+ ${turbine.version}
+
+
+ com.netflix.rxjava
+ rxjava-core
+
+
+ org.slf4j
+ slf4j-simple
+
+
+
+
+
+
+
+
+ org.springframework.cloud
+ spring-cloud-netflix-core
+
+
+ org.springframework.boot
+ spring-boot-starter-amqp
+
+
+ org.springframework.boot
+ spring-boot-starter-integration
+
+
+ org.springframework.integration
+ spring-integration-amqp
+ true
+
+
+ org.springframework.integration
+ spring-integration-java-dsl
+ true
+
+
+ com.fasterxml.jackson.core
+ jackson-databind
+ true
+
+
+ com.netflix.turbine
+ turbine-core
+ true
+
+
+ io.reactivex
+ rxjava
+ true
+
+
+ org.projectlombok
+ lombok
+ compile
+ true
+
+
+
+
diff --git a/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/netflix/turbine/amqp/Aggregator.java b/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/netflix/turbine/amqp/Aggregator.java
new file mode 100644
index 00000000..505961d2
--- /dev/null
+++ b/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/netflix/turbine/amqp/Aggregator.java
@@ -0,0 +1,65 @@
+package org.springframework.netflix.turbine.amqp;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+
+import lombok.extern.slf4j.Slf4j;
+
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.integration.annotation.MessageEndpoint;
+import org.springframework.integration.annotation.ServiceActivator;
+import org.springframework.util.StringUtils;
+
+import rx.subjects.PublishSubject;
+
+import java.io.IOException;
+import java.util.Map;
+
+/**
+ * @author Spencer Gibb
+ */
+@MessageEndpoint
+@Slf4j
+public class Aggregator {
+
+ @Autowired
+ private ObjectMapper objectMapper;
+
+ @Autowired
+ private PublishSubject