From a912000f0fdf20039e1df1f7edff2d4155e23ef0 Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Mon, 29 Jul 2019 12:32:38 +0200 Subject: [PATCH] Added option to send spans to Zipkin over ActiveMQ fixes gh-1380 --- spring-cloud-sleuth-zipkin/pom.xml | 15 +++++ .../ZipkinActiveMqSenderConfiguration.java | 61 +++++++++++++++++++ ...pkinSenderConfigurationImportSelector.java | 1 + .../sender/ZipkinSenderProperties.java | 5 ++ ...itional-spring-configuration-metadata.json | 28 +++++++++ .../zipkin2/ZipkinAutoConfigurationTests.java | 20 ++++++ 6 files changed, 130 insertions(+) create mode 100644 spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/sender/ZipkinActiveMqSenderConfiguration.java create mode 100644 spring-cloud-sleuth-zipkin/src/main/resources/META-INF/additional-spring-configuration-metadata.json diff --git a/spring-cloud-sleuth-zipkin/pom.xml b/spring-cloud-sleuth-zipkin/pom.xml index 52518583a..ba49fe9a7 100644 --- a/spring-cloud-sleuth-zipkin/pom.xml +++ b/spring-cloud-sleuth-zipkin/pom.xml @@ -89,6 +89,21 @@ + + io.zipkin.reporter2 + zipkin-sender-activemq-client + + + org.apache.activemq + activemq-client + + + + + org.apache.activemq + activemq-client + true + org.springframework.kafka spring-kafka diff --git a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/sender/ZipkinActiveMqSenderConfiguration.java b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/sender/ZipkinActiveMqSenderConfiguration.java new file mode 100644 index 000000000..f33cfc325 --- /dev/null +++ b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/sender/ZipkinActiveMqSenderConfiguration.java @@ -0,0 +1,61 @@ +/* + * Copyright 2013-2019 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.zipkin2.sender; + +import org.apache.activemq.ActiveMQConnectionFactory; +import zipkin2.reporter.Sender; +import zipkin2.reporter.activemq.ActiveMQSender; + +import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.autoconfigure.AutoConfigureAfter; +import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.boot.autoconfigure.jms.activemq.ActiveMQAutoConfiguration; +import org.springframework.cloud.sleuth.zipkin2.ZipkinAutoConfiguration; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Conditional; +import org.springframework.context.annotation.Configuration; + +@Configuration +@ConditionalOnClass(ActiveMQConnectionFactory.class) +@ConditionalOnMissingBean(name = ZipkinAutoConfiguration.SENDER_BEAN_NAME) +@Conditional(ZipkinSenderCondition.class) +@ConditionalOnProperty(value = "spring.zipkin.sender.type", havingValue = "activemq") +@AutoConfigureAfter(ActiveMQAutoConfiguration.class) +class ZipkinActiveMqSenderConfiguration { + + @Configuration + @ConditionalOnBean(ActiveMQConnectionFactory.class) + static class ZipkinActiveMqSenderBeanConfiguration { + + @Value("${spring.zipkin.activemq.queue:zipkin}") + private String queue; + + @Value("${spring.zipkin.activemq.message-max-bytes:100000}") + private int messageMaxBytes; + + @Bean(ZipkinAutoConfiguration.SENDER_BEAN_NAME) + Sender activeMqSender(ActiveMQConnectionFactory factory) { + return ActiveMQSender.newBuilder().connectionFactory(factory) + .messageMaxBytes(this.messageMaxBytes).queue(this.queue).build(); + } + + } + +} diff --git a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/sender/ZipkinSenderConfigurationImportSelector.java b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/sender/ZipkinSenderConfigurationImportSelector.java index e0a79c638..41857c00f 100644 --- a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/sender/ZipkinSenderConfigurationImportSelector.java +++ b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/sender/ZipkinSenderConfigurationImportSelector.java @@ -36,6 +36,7 @@ public class ZipkinSenderConfigurationImportSelector implements ImportSelector { static { // Mappings in descending priority (highest is last) Map mappings = new LinkedHashMap<>(); + mappings.put("activemq", ZipkinActiveMqSenderConfiguration.class.getName()); mappings.put("rabbit", ZipkinRabbitSenderConfiguration.class.getName()); mappings.put("kafka", ZipkinKafkaSenderConfiguration.class.getName()); mappings.put("web", ZipkinRestTemplateSenderConfiguration.class.getName()); diff --git a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/sender/ZipkinSenderProperties.java b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/sender/ZipkinSenderProperties.java index 067b85908..7037a36b3 100644 --- a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/sender/ZipkinSenderProperties.java +++ b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/sender/ZipkinSenderProperties.java @@ -45,6 +45,11 @@ public class ZipkinSenderProperties { */ public enum SenderType { + /** + * ActiveMQ sender. + */ + ACTIVEMQ, + /** * RabbitMQ sender. */ diff --git a/spring-cloud-sleuth-zipkin/src/main/resources/META-INF/additional-spring-configuration-metadata.json b/spring-cloud-sleuth-zipkin/src/main/resources/META-INF/additional-spring-configuration-metadata.json new file mode 100644 index 000000000..9527605f0 --- /dev/null +++ b/spring-cloud-sleuth-zipkin/src/main/resources/META-INF/additional-spring-configuration-metadata.json @@ -0,0 +1,28 @@ +{ + "properties": [ + { + "name": "spring.zipkin.kafka.topic", + "type": "java.lang.String", + "description": "Name of the Kafka topic where spans should be sent to Zipkin.", + "defaultValue": "zipkin" + }, + { + "name": "spring.zipkin.rabbitmq.queue", + "type": "java.lang.String", + "description": "Name of the RabbitMQ queue where spans should be sent to Zipkin.", + "defaultValue": "zipkin" + }, + { + "name": "spring.zipkin.activemq.queue", + "type": "java.lang.String", + "description": "Name of the ActiveMQ queue where spans should be sent to Zipkin.", + "defaultValue": "zipkin" + }, + { + "name": "spring.zipkin.activemq.message-max-bytes", + "type": "java.lang.String", + "description": "Maximum number of bytes for a given message with spans sent to Zipkin over ActiveMQ.", + "defaultValue": 100000 + } + ] +} diff --git a/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin2/ZipkinAutoConfigurationTests.java b/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin2/ZipkinAutoConfigurationTests.java index 2d290a316..acae31d10 100644 --- a/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin2/ZipkinAutoConfigurationTests.java +++ b/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin2/ZipkinAutoConfigurationTests.java @@ -36,11 +36,13 @@ import zipkin2.codec.Encoding; import zipkin2.reporter.AsyncReporter; import zipkin2.reporter.Reporter; import zipkin2.reporter.Sender; +import zipkin2.reporter.activemq.ActiveMQSender; import zipkin2.reporter.amqp.RabbitMQSender; import zipkin2.reporter.kafka.KafkaSender; import org.springframework.boot.autoconfigure.amqp.RabbitAutoConfiguration; import org.springframework.boot.autoconfigure.context.PropertyPlaceholderAutoConfiguration; +import org.springframework.boot.autoconfigure.jms.activemq.ActiveMQAutoConfiguration; import org.springframework.boot.autoconfigure.kafka.KafkaAutoConfiguration; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.sleuth.autoconfig.TraceAutoConfiguration; @@ -155,6 +157,24 @@ public class ZipkinAutoConfigurationTests { this.context.close(); } + @Test + public void overrideActiveMqQueue() throws Exception { + this.context = new AnnotationConfigApplicationContext(); + environment().setProperty("spring.jms.cache.enabled", "false"); + environment().setProperty("spring.zipkin.activemq.queue", "zipkin2"); + environment().setProperty("spring.zipkin.activemq.message-max-bytes", "50"); + environment().setProperty("spring.zipkin.sender.type", "activemq"); + this.context.register(PropertyPlaceholderAutoConfiguration.class, + ActiveMQAutoConfiguration.class, ZipkinAutoConfiguration.class, + TraceAutoConfiguration.class, + ZipkinBackwardsCompatibilityAutoConfiguration.class); + this.context.refresh(); + + then(this.context.getBean(Sender.class)).isInstanceOf(ActiveMQSender.class); + + this.context.close(); + } + @Test public void canOverrideBySender() throws Exception { this.context = new AnnotationConfigApplicationContext();