diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/config/CollectorConfig.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/config/CollectorConfig.java new file mode 100644 index 000000000..54d9b3f0e --- /dev/null +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/config/CollectorConfig.java @@ -0,0 +1,15 @@ +package com.a.eye.skywalking.collector.cluster.config; + +/** + * Created by pengys5 on 2017/2/22 0022. + */ +public class CollectorConfig { + + public static final String appname = "CollectorSystem"; + + public static class Collector { + public static String hostname = "127.0.0.1"; + public static String port = "2551"; + public static String cluster = "127.0.0.1:2551"; + } +} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/config/CollectorConfigInitializer.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/config/CollectorConfigInitializer.java new file mode 100644 index 000000000..73c0dad75 --- /dev/null +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/config/CollectorConfigInitializer.java @@ -0,0 +1,44 @@ +package com.a.eye.skywalking.collector.cluster.config; + +import com.a.eye.skywalking.api.conf.Config; +import com.a.eye.skywalking.api.logging.api.ILog; +import com.a.eye.skywalking.api.logging.api.LogManager; +import com.a.eye.skywalking.api.util.ConfigInitializer; +import com.a.eye.skywalking.api.util.StringUtil; + +import java.io.InputStream; +import java.util.Properties; + +/** + * Created by pengys5 on 2017/2/22 0022. + */ +public class CollectorConfigInitializer { + + private static ILog logger = LogManager.getLogger(CollectorConfigInitializer.class); + + public static void initialize() { + InputStream configFileStream = CollectorConfigInitializer.class.getResourceAsStream("/collector.config"); + + if (configFileStream == null) { + logger.info("Not provide sky-walking certification documents, sky-walking api run in default config."); + } else { + try { + Properties properties = new Properties(); + properties.load(configFileStream); + ConfigInitializer.initialize(properties, CollectorConfig.class); + } catch (Exception e) { + logger.error("Failed to read the config file, sky-walking api run in default config.", e); + } + } + + if (!StringUtil.isEmpty(System.getProperty("collector.hostname"))) { + CollectorConfig.Collector.hostname = System.getProperty("collector.hostname"); + } + if (!StringUtil.isEmpty(System.getProperty("collector.port"))) { + CollectorConfig.Collector.port = System.getProperty("collector.port"); + } + if (!StringUtil.isEmpty(System.getProperty("collector.cluster"))) { + CollectorConfig.Collector.cluster = System.getProperty("collector.cluster"); + } + } +} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/producer/TraceProducerApp.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/producer/TraceProducerApp.java index e9452121d..86be204ca 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/producer/TraceProducerApp.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/producer/TraceProducerApp.java @@ -1,37 +1,38 @@ package com.a.eye.skywalking.collector.cluster.producer; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicInteger; - -import com.a.eye.skywalking.collector.cluster.Const; -import com.a.eye.skywalking.collector.cluster.consumer.TraceConsumerActor; -import com.a.eye.skywalking.collector.cluster.manager.ActorManagerActor; -import com.a.eye.skywalking.sniffer.mock.trace.TraceSegmentBuilderFactory; -import com.a.eye.skywalking.trace.TraceSegment; -import com.typesafe.config.Config; -import com.typesafe.config.ConfigFactory; - -import com.a.eye.skywalking.collector.cluster.message.TraceMessages.TransformationJob; -import scala.concurrent.ExecutionContext; -import scala.concurrent.duration.Duration; -import scala.concurrent.duration.FiniteDuration; import akka.actor.ActorRef; import akka.actor.ActorSystem; import akka.actor.Props; import akka.dispatch.OnSuccess; import akka.util.Timeout; +import com.a.eye.skywalking.collector.cluster.Const; +import com.a.eye.skywalking.collector.cluster.config.CollectorConfig; +import com.a.eye.skywalking.collector.cluster.config.CollectorConfigInitializer; +import com.a.eye.skywalking.collector.cluster.manager.ActorManagerActor; +import com.a.eye.skywalking.collector.cluster.message.TraceMessages.TransformationJob; +import com.typesafe.config.Config; +import com.typesafe.config.ConfigFactory; +import scala.concurrent.ExecutionContext; +import scala.concurrent.duration.Duration; +import scala.concurrent.duration.FiniteDuration; + +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import static akka.pattern.Patterns.ask; +/** + * {@link TraceProducerApp} is a producer for trace agent to send @link TraceSegment. + *
+ * Created by pengys5 on 2017/2/17. + */ public class TraceProducerApp { public static void main(String[] args) { // Override the configuration of the port when specified as program argument - final String port = args.length > 0 ? args[0] : "2552"; - final Config config = ConfigFactory.parseString("akka.remote.netty.tcp.port=" + port). - withFallback(ConfigFactory.load()); + final Config config = TraceProducerApp.buildConfig(); - ActorSystem system = ActorSystem.create("ClusterSystem", config); + ActorSystem system = ActorSystem.create(CollectorConfig.appname, config); system.actorOf(Props.create(ActorManagerActor.class), Const.Actor_Manager_Role); final ActorRef frontend = system.actorOf(Props.create(TraceProducerActor.class), Const.Trace_Producer_Role); @@ -39,19 +40,44 @@ public class TraceProducerApp { final Timeout timeout = new Timeout(Duration.create(5, TimeUnit.SECONDS)); final ExecutionContext ec = system.dispatcher(); final AtomicInteger counter = new AtomicInteger(); - system.scheduler().schedule(interval, interval, new Runnable() { - public void run() { -// TraceSegment traceSegment = TraceSegmentBuilderFactory.INSTANCE.singleTomcat200Trace(); - ask(frontend, - new TransformationJob("hello-" + counter.incrementAndGet(), null), - timeout).onSuccess(new OnSuccess