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 bdd0f340..304c246d 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 @@ -26,7 +26,7 @@ public class HystrixStreamAutoConfiguration { @ConditionalOnExpression("${hystrix.stream.amqp.enabled:true}") @IntegrationComponentScan(basePackageClasses = HystrixStreamChannel.class) @EnableScheduling - protected static class HystrixStreamBusAutoConfiguration { + protected static class HystrixStreamAmqpAutoConfiguration { @Autowired private AmqpTemplate amqpTemplate; diff --git a/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/netflix/turbine/amqp/EnableTurbineAmqp.java b/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/netflix/turbine/amqp/EnableTurbineAmqp.java index 66d86545..d05ffa1c 100644 --- a/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/netflix/turbine/amqp/EnableTurbineAmqp.java +++ b/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/netflix/turbine/amqp/EnableTurbineAmqp.java @@ -5,7 +5,7 @@ import org.springframework.context.annotation.Import; import java.lang.annotation.*; /** - * Run the RxNetty based Spring Cloud Bus Turbine server. + * Run the RxNetty based Spring Cloud Turbine AMQP server. * Based on Netflix Turbine 2 * * @author Spencer Gibb diff --git a/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/netflix/turbine/amqp/TurbineAmqpAutoConfiguration.java b/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/netflix/turbine/amqp/TurbineAmqpAutoConfiguration.java index 48327831..eb33a35d 100644 --- a/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/netflix/turbine/amqp/TurbineAmqpAutoConfiguration.java +++ b/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/netflix/turbine/amqp/TurbineAmqpAutoConfiguration.java @@ -27,7 +27,7 @@ import org.springframework.integration.dsl.amqp.Amqp; public class TurbineAmqpAutoConfiguration { @Configuration - @ConditionalOnExpression("${hystrix.stream.bus.turbine.enabled:true}") + @ConditionalOnExpression("${turbine.amqp.enabled:true}") protected static class HystrixStreamAggregatorAutoConfiguration { @Autowired @@ -41,7 +41,7 @@ public class TurbineAmqpAutoConfiguration { } @Bean - protected Binding localCloudBusQueueBinding() { + protected Binding localTurbineAmqpQueueBinding() { return BindingBuilder.bind(hystrixStreamQueue()).to(hystrixStreamExchange()).with(""); } diff --git a/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/netflix/turbine/amqp/TurbineAmqpConfiguration.java b/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/netflix/turbine/amqp/TurbineAmqpConfiguration.java index b1353eef..b9b1a777 100644 --- a/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/netflix/turbine/amqp/TurbineAmqpConfiguration.java +++ b/spring-cloud-netflix-turbine-amqp/src/main/java/org/springframework/netflix/turbine/amqp/TurbineAmqpConfiguration.java @@ -24,7 +24,7 @@ import static io.reactivex.netty.pipeline.PipelineConfigurators.sseServerConfigu */ @Configuration @Slf4j -@ConfigurationProperties("bus.turbine") +@ConfigurationProperties("turbine.amqp") public class TurbineAmqpConfiguration implements SmartLifecycle { private boolean running = false; @@ -41,15 +41,15 @@ public class TurbineAmqpConfiguration implements SmartLifecycle { // multicast so multiple concurrent subscribers get the same stream Observable> publishedStreams = StreamAggregator.aggregateGroupedStreams(hystrixSubject() .groupBy(data -> InstanceKey.create((String) data.get("instanceId")))) - .doOnUnsubscribe(() -> log.info("BusTurbine => Unsubscribing aggregation.")) - .doOnSubscribe(() -> log.info("BusTurbine => Starting aggregation")) + .doOnUnsubscribe(() -> log.info("AmqpTurbine => Unsubscribing aggregation.")) + .doOnSubscribe(() -> log.info("AmqpTurbine => Starting aggregation")) .flatMap(o -> o).publish().refCount(); HttpServer httpServer = RxNetty.createHttpServer(port, (request, response) -> { - log.info("BusTurbine => SSE Request Received"); + log.info("AmqpTurbine => SSE Request Received"); response.getHeaders().setHeader("Content-Type", "text/event-stream"); return publishedStreams - .doOnUnsubscribe(() -> log.info("BusTurbine => Unsubscribing RxNetty server connection")) + .doOnUnsubscribe(() -> log.info("AmqpTurbine => Unsubscribing RxNetty server connection")) .flatMap(data -> response.writeAndFlush(new ServerSentEvent(null, null, JsonUtility.mapToJson(data)))); }, sseServerConfigurator()); return httpServer; diff --git a/spring-cloud-netflix-turbine-amqp/src/test/java/org/springframework/netflix/turbine/amqp/AggregatorTest.java b/spring-cloud-netflix-turbine-amqp/src/test/java/org/springframework/netflix/turbine/amqp/AggregatorTest.java index 019d18b2..3e6f05b6 100644 --- a/spring-cloud-netflix-turbine-amqp/src/test/java/org/springframework/netflix/turbine/amqp/AggregatorTest.java +++ b/spring-cloud-netflix-turbine-amqp/src/test/java/org/springframework/netflix/turbine/amqp/AggregatorTest.java @@ -17,7 +17,7 @@ import static org.springframework.netflix.turbine.amqp.Aggregator.getPayloadData public class AggregatorTest { - public static final String STREAM_ALL = "hystrixbus"; + public static final String STREAM_ALL = "hystrixamqp"; public static void main(String[] args) { getHystrixStreamFromFile(STREAM_ALL, 1).flatMap(commandGroup -> commandGroup.take(50)).take(50).toBlocking().forEach(s -> System.out.println("s: " + s)); diff --git a/spring-cloud-netflix-turbine-amqp/src/test/resources/org/springframework/netflix/turbine/amqp/hystrixbus.stream b/spring-cloud-netflix-turbine-amqp/src/test/resources/org/springframework/netflix/turbine/amqp/hystrixamqp.stream similarity index 100% rename from spring-cloud-netflix-turbine-amqp/src/test/resources/org/springframework/netflix/turbine/amqp/hystrixbus.stream rename to spring-cloud-netflix-turbine-amqp/src/test/resources/org/springframework/netflix/turbine/amqp/hystrixamqp.stream