From 911de3fe4a500ca01ea69f29192bbad3c0909d23 Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Sun, 6 Aug 2017 21:54:55 +0800 Subject: [PATCH] 1. Provide jetty http interface for application and instance register. 2. Provide jetty http interface for receive trace segment. 3. Test segment receive trace segment success. segment, node reference, node reference sum, node component, node mapping, global trace data are all correct. --- .../jetty/AgentRegisterJettyConfig.java | 10 ++ .../jetty/AgentRegisterJettyConfigParser.java | 36 +++++++ .../jetty/AgentRegisterJettyDataListener.java | 21 ++++ .../jetty/AgentRegisterJettyModuleDefine.java | 53 ++++++++++ .../AgentRegisterJettyModuleRegistration.java | 13 +++ .../ApplicationRegisterServletHandler.java | 52 ++++++++++ .../InstanceDiscoveryServletHandler.java | 53 ++++++++++ .../resources/META-INF/defines/module.define | 3 +- .../handler/AgentStreamGRPCServerHandler.java | 6 +- .../AgentStreamJettyServerHandler.java | 2 +- .../jetty/AgentStreamJettyModuleDefine.java | 6 +- .../handler/TraceSegmentServiceHandler.java | 34 ------- .../handler/TraceSegmentServletHandler.java | 54 ++++++++++ .../reader/KeyWithStringValueJsonReader.java | 36 +++++++ .../jetty/handler/reader/LogJsonReader.java | 37 +++++++ .../handler/reader/ReferenceJsonReader.java | 70 +++++++++++++ .../handler/reader/SegmentJsonReader.java | 69 +++++++++++++ .../jetty/handler/reader/SpanJsonReader.java | 99 +++++++++++++++++++ .../handler/reader/StreamJsonReader.java | 11 +++ .../jetty/handler/reader/TraceSegment.java | 35 +++++++ .../reader/TraceSegmentJsonReader.java | 54 ++++++++++ .../handler/reader/UniqueIdJsonReader.java | 22 +++++ .../collector/agentstream/worker/Const.java | 1 + .../worker/cache/InstanceCache.java | 25 +++++ .../component/NodeComponentSpanListener.java | 4 +- .../node/mapping/NodeMappingSpanListener.java | 4 +- .../reference/NodeRefSpanListener.java | 23 +++-- .../summary/NodeRefSumSpanListener.java | 27 ++++- .../ApplicationRegisterSerialWorker.java | 8 +- .../register/instance/dao/IInstanceDAO.java | 2 + .../register/instance/dao/InstanceEsDAO.java | 17 +++- .../register/instance/dao/InstanceH2DAO.java | 4 + .../worker/storage/PersistenceTimer.java | 32 +++--- .../agentstream/HttpClientTools.java | 77 +++++++++++++++ .../TraceSegmentJsonReaderTestCase.java | 28 ++++++ .../agentstream/mock/JsonFileReader.java | 24 +++++ .../agentstream/mock/SegmentPost.java | 39 ++++++++ .../segment/normal/application-register.json | 4 + .../json/segment/normal/dubbox-consumer.json | 71 +++++++++++++ .../json/segment/normal/dubbox-provider.json | 66 +++++++++++++ .../segment/normal/instance-register.json | 5 + .../elasticsearch/ElasticSearchClient.java | 1 - .../collector/server/jetty/JettyHandler.java | 22 ++--- .../storage/elasticsearch/dao/BatchEsDAO.java | 12 ++- .../jetty/handler/SegmentTopGetHandler.java | 2 +- .../ui/jetty/handler/SpanGetHandler.java | 2 +- .../ui/jetty/handler/TraceDagGetHandler.java | 2 +- .../jetty/handler/TraceStackGetHandler.java | 2 +- .../jetty/handler/UIJettyServerHandler.java | 2 +- 49 files changed, 1183 insertions(+), 99 deletions(-) create mode 100644 apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyConfig.java create mode 100644 apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyConfigParser.java create mode 100644 apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyDataListener.java create mode 100644 apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyModuleDefine.java create mode 100644 apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyModuleRegistration.java create mode 100644 apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/handler/ApplicationRegisterServletHandler.java create mode 100644 apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/handler/InstanceDiscoveryServletHandler.java delete mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/TraceSegmentServiceHandler.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/TraceSegmentServletHandler.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/KeyWithStringValueJsonReader.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/LogJsonReader.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/ReferenceJsonReader.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/SegmentJsonReader.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/SpanJsonReader.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/StreamJsonReader.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/TraceSegment.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/TraceSegmentJsonReader.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/UniqueIdJsonReader.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java create mode 100644 apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/HttpClientTools.java create mode 100644 apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/TraceSegmentJsonReaderTestCase.java create mode 100644 apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/JsonFileReader.java create mode 100644 apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java create mode 100644 apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/application-register.json create mode 100644 apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/dubbox-consumer.json create mode 100644 apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/dubbox-provider.json create mode 100644 apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/instance-register.json diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyConfig.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyConfig.java new file mode 100644 index 000000000..cf1e45207 --- /dev/null +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyConfig.java @@ -0,0 +1,10 @@ +package org.skywalking.apm.collector.agentregister.jetty; + +/** + * @author pengys5 + */ +public class AgentRegisterJettyConfig { + public static String HOST; + public static int PORT; + public static String CONTEXT_PATH; +} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyConfigParser.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyConfigParser.java new file mode 100644 index 000000000..65a10bebe --- /dev/null +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyConfigParser.java @@ -0,0 +1,36 @@ +package org.skywalking.apm.collector.agentregister.jetty; + +import java.util.Map; +import org.skywalking.apm.collector.core.config.ConfigParseException; +import org.skywalking.apm.collector.core.module.ModuleConfigParser; +import org.skywalking.apm.collector.core.util.ObjectUtils; +import org.skywalking.apm.collector.core.util.StringUtils; + +/** + * @author pengys5 + */ +public class AgentRegisterJettyConfigParser implements ModuleConfigParser { + + private static final String HOST = "host"; + private static final String PORT = "port"; + public static final String CONTEXT_PATH = "contextPath"; + + @Override public void parse(Map config) throws ConfigParseException { + AgentRegisterJettyConfig.CONTEXT_PATH = "/"; + + if (ObjectUtils.isEmpty(config) || StringUtils.isEmpty(config.get(HOST))) { + AgentRegisterJettyConfig.HOST = "localhost"; + } else { + AgentRegisterJettyConfig.HOST = (String)config.get(HOST); + } + + if (ObjectUtils.isEmpty(config) || StringUtils.isEmpty(config.get(PORT))) { + AgentRegisterJettyConfig.PORT = 12800; + } else { + AgentRegisterJettyConfig.PORT = (Integer)config.get(PORT); + } + if (ObjectUtils.isNotEmpty(config) && StringUtils.isNotEmpty(config.get(CONTEXT_PATH))) { + AgentRegisterJettyConfig.CONTEXT_PATH = (String)config.get(CONTEXT_PATH); + } + } +} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyDataListener.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyDataListener.java new file mode 100644 index 000000000..218e0f42d --- /dev/null +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyDataListener.java @@ -0,0 +1,21 @@ +package org.skywalking.apm.collector.agentregister.jetty; + +import org.skywalking.apm.collector.agentregister.AgentRegisterModuleGroupDefine; +import org.skywalking.apm.collector.agentregister.grpc.AgentRegisterGRPCModuleDefine; +import org.skywalking.apm.collector.cluster.ClusterModuleDefine; +import org.skywalking.apm.collector.core.cluster.ClusterDataListener; + +/** + * @author pengys5 + */ +public class AgentRegisterJettyDataListener extends ClusterDataListener { + + public static final String PATH = ClusterModuleDefine.BASE_CATALOG + "." + AgentRegisterModuleGroupDefine.GROUP_NAME + "." + AgentRegisterJettyModuleDefine.MODULE_NAME; + + @Override public String path() { + return PATH; + } + + @Override public void addressChangedNotify() { + } +} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyModuleDefine.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyModuleDefine.java new file mode 100644 index 000000000..413149d71 --- /dev/null +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyModuleDefine.java @@ -0,0 +1,53 @@ +package org.skywalking.apm.collector.agentregister.jetty; + +import java.util.LinkedList; +import java.util.List; +import org.skywalking.apm.collector.agentregister.AgentRegisterModuleDefine; +import org.skywalking.apm.collector.agentregister.AgentRegisterModuleGroupDefine; +import org.skywalking.apm.collector.agentregister.jetty.handler.ApplicationRegisterServletHandler; +import org.skywalking.apm.collector.agentregister.jetty.handler.InstanceDiscoveryServletHandler; +import org.skywalking.apm.collector.core.cluster.ClusterDataListener; +import org.skywalking.apm.collector.core.framework.Handler; +import org.skywalking.apm.collector.core.module.ModuleConfigParser; +import org.skywalking.apm.collector.core.module.ModuleRegistration; +import org.skywalking.apm.collector.core.server.Server; +import org.skywalking.apm.collector.server.jetty.JettyServer; + +/** + * @author pengys5 + */ +public class AgentRegisterJettyModuleDefine extends AgentRegisterModuleDefine { + + public static final String MODULE_NAME = "jetty"; + + @Override protected String group() { + return AgentRegisterModuleGroupDefine.GROUP_NAME; + } + + @Override public String name() { + return MODULE_NAME; + } + + @Override protected ModuleConfigParser configParser() { + return new AgentRegisterJettyConfigParser(); + } + + @Override protected Server server() { + return new JettyServer(AgentRegisterJettyConfig.HOST, AgentRegisterJettyConfig.PORT, AgentRegisterJettyConfig.CONTEXT_PATH); + } + + @Override protected ModuleRegistration registration() { + return new AgentRegisterJettyModuleRegistration(); + } + + @Override public ClusterDataListener listener() { + return new AgentRegisterJettyDataListener(); + } + + @Override public List handlerList() { + List handlers = new LinkedList<>(); + handlers.add(new ApplicationRegisterServletHandler()); + handlers.add(new InstanceDiscoveryServletHandler()); + return handlers; + } +} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyModuleRegistration.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyModuleRegistration.java new file mode 100644 index 000000000..d6e2d260a --- /dev/null +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyModuleRegistration.java @@ -0,0 +1,13 @@ +package org.skywalking.apm.collector.agentregister.jetty; + +import org.skywalking.apm.collector.core.module.ModuleRegistration; + +/** + * @author pengys5 + */ +public class AgentRegisterJettyModuleRegistration extends ModuleRegistration { + + @Override public Value buildValue() { + return new Value(AgentRegisterJettyConfig.HOST, AgentRegisterJettyConfig.PORT, AgentRegisterJettyConfig.CONTEXT_PATH); + } +} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/handler/ApplicationRegisterServletHandler.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/handler/ApplicationRegisterServletHandler.java new file mode 100644 index 000000000..7c7f2a99b --- /dev/null +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/handler/ApplicationRegisterServletHandler.java @@ -0,0 +1,52 @@ +package org.skywalking.apm.collector.agentregister.jetty.handler; + +import com.google.gson.Gson; +import com.google.gson.JsonArray; +import com.google.gson.JsonElement; +import com.google.gson.JsonObject; +import java.io.IOException; +import javax.servlet.http.HttpServletRequest; +import org.skywalking.apm.collector.agentregister.application.ApplicationIDService; +import org.skywalking.apm.collector.server.jetty.ArgumentsParseException; +import org.skywalking.apm.collector.server.jetty.JettyHandler; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class ApplicationRegisterServletHandler extends JettyHandler { + + private final Logger logger = LoggerFactory.getLogger(ApplicationRegisterServletHandler.class); + + private ApplicationIDService applicationIDService = new ApplicationIDService(); + private Gson gson = new Gson(); + private static final String APPLICATION_CODE = "c"; + private static final String APPLICATION_ID = "i"; + + @Override public String pathSpec() { + return "/application/register"; + } + + @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { + throw new UnsupportedOperationException(); + } + + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { + JsonArray responseArray = new JsonArray(); + try { + JsonArray applicationCodes = gson.fromJson(req.getReader(), JsonArray.class); + for (int i = 0; i < applicationCodes.size(); i++) { + String applicationCode = applicationCodes.get(i).getAsString(); + int applicationId = applicationIDService.getOrCreate(applicationCode); + JsonObject mapping = new JsonObject(); + mapping.addProperty(APPLICATION_CODE, applicationCode); + mapping.addProperty(APPLICATION_ID, applicationId); + responseArray.add(mapping); + } + } catch (IOException e) { + logger.error(e.getMessage(), e); + } + return responseArray; + } +} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/handler/InstanceDiscoveryServletHandler.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/handler/InstanceDiscoveryServletHandler.java new file mode 100644 index 000000000..ffddfb94d --- /dev/null +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/handler/InstanceDiscoveryServletHandler.java @@ -0,0 +1,53 @@ +package org.skywalking.apm.collector.agentregister.jetty.handler; + +import com.google.gson.Gson; +import com.google.gson.JsonElement; +import com.google.gson.JsonObject; +import java.io.IOException; +import javax.servlet.http.HttpServletRequest; +import org.skywalking.apm.collector.agentregister.instance.InstanceIDService; +import org.skywalking.apm.collector.server.jetty.ArgumentsParseException; +import org.skywalking.apm.collector.server.jetty.JettyHandler; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class InstanceDiscoveryServletHandler extends JettyHandler { + + private final Logger logger = LoggerFactory.getLogger(InstanceDiscoveryServletHandler.class); + + private InstanceIDService instanceIDService = new InstanceIDService(); + private Gson gson = new Gson(); + + private static final String APPLICATION_ID = "ai"; + private static final String AGENT_UUID = "au"; + private static final String REGISTER_TIME = "rt"; + private static final String INSTANCE_ID = "ii"; + + @Override public String pathSpec() { + return "/instance/register"; + } + + @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { + throw new UnsupportedOperationException(); + } + + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { + JsonObject responseJson = new JsonObject(); + try { + JsonObject instance = gson.fromJson(req.getReader(), JsonObject.class); + int applicationId = instance.get(APPLICATION_ID).getAsInt(); + String agentUUID = instance.get(AGENT_UUID).getAsString(); + long registerTime = instance.get(REGISTER_TIME).getAsLong(); + + int instanceId = instanceIDService.getOrCreate(applicationId, agentUUID, registerTime); + responseJson.addProperty(APPLICATION_ID, applicationId); + responseJson.addProperty(INSTANCE_ID, instanceId); + } catch (IOException e) { + logger.error(e.getMessage(), e); + } + return responseJson; + } +} diff --git a/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/module.define b/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/module.define index 70f99f0a2..e826d4160 100644 --- a/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/module.define +++ b/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/module.define @@ -1 +1,2 @@ -org.skywalking.apm.collector.agentregister.grpc.AgentRegisterGRPCModuleDefine \ No newline at end of file +org.skywalking.apm.collector.agentregister.grpc.AgentRegisterGRPCModuleDefine +org.skywalking.apm.collector.agentregister.jetty.AgentRegisterJettyModuleDefine \ No newline at end of file diff --git a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamGRPCServerHandler.java b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamGRPCServerHandler.java index db9b9362a..f1cd49461 100644 --- a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamGRPCServerHandler.java +++ b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamGRPCServerHandler.java @@ -25,13 +25,11 @@ public class AgentStreamGRPCServerHandler extends JettyHandler { ClusterModuleRegistrationReader reader = ((ClusterModuleContext)CollectorContextHelper.INSTANCE.getContext(ClusterModuleGroupDefine.GROUP_NAME)).getReader(); List servers = reader.read(AgentStreamGRPCDataListener.PATH); JsonArray serverArray = new JsonArray(); - servers.forEach(server -> { - serverArray.add(server); - }); + servers.forEach(server -> serverArray.add(server)); return serverArray; } - @Override protected void doPost(HttpServletRequest req) throws ArgumentsParseException { + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { throw new UnsupportedOperationException(); } } diff --git a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamJettyServerHandler.java b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamJettyServerHandler.java index a463de01d..bfcecdbf7 100644 --- a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamJettyServerHandler.java +++ b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamJettyServerHandler.java @@ -31,7 +31,7 @@ public class AgentStreamJettyServerHandler extends JettyHandler { return serverArray; } - @Override protected void doPost(HttpServletRequest req) throws ArgumentsParseException { + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { throw new UnsupportedOperationException(); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyModuleDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyModuleDefine.java index d20996ebd..87428c339 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyModuleDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyModuleDefine.java @@ -1,8 +1,10 @@ package org.skywalking.apm.collector.agentstream.jetty; +import java.util.LinkedList; import java.util.List; import org.skywalking.apm.collector.agentstream.AgentStreamModuleDefine; import org.skywalking.apm.collector.agentstream.AgentStreamModuleGroupDefine; +import org.skywalking.apm.collector.agentstream.jetty.handler.TraceSegmentServletHandler; import org.skywalking.apm.collector.core.cluster.ClusterDataListener; import org.skywalking.apm.collector.core.framework.Handler; import org.skywalking.apm.collector.core.module.ModuleConfigParser; @@ -42,6 +44,8 @@ public class AgentStreamJettyModuleDefine extends AgentStreamModuleDefine { } @Override public List handlerList() { - return null; + List handlers = new LinkedList<>(); + handlers.add(new TraceSegmentServletHandler()); + return handlers; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/TraceSegmentServiceHandler.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/TraceSegmentServiceHandler.java deleted file mode 100644 index a92102319..000000000 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/TraceSegmentServiceHandler.java +++ /dev/null @@ -1,34 +0,0 @@ -package org.skywalking.apm.collector.agentstream.jetty.handler; - -import com.google.gson.JsonElement; -import java.io.BufferedReader; -import java.io.IOException; -import javax.servlet.http.HttpServletRequest; -import org.skywalking.apm.collector.server.jetty.ArgumentsParseException; -import org.skywalking.apm.collector.server.jetty.JettyHandler; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -/** - * @author pengys5 - */ -public class TraceSegmentServiceHandler extends JettyHandler { - - private final Logger logger = LoggerFactory.getLogger(TraceSegmentServiceHandler.class); - - @Override public String pathSpec() { - return null; - } - - @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { - throw new UnsupportedOperationException(); - } - - @Override protected void doPost(HttpServletRequest req) throws ArgumentsParseException { - try { - BufferedReader reader = req.getReader(); - } catch (IOException e) { - logger.error(e.getMessage(), e); - } - } -} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/TraceSegmentServletHandler.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/TraceSegmentServletHandler.java new file mode 100644 index 000000000..7622aac48 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/TraceSegmentServletHandler.java @@ -0,0 +1,54 @@ +package org.skywalking.apm.collector.agentstream.jetty.handler; + +import com.google.gson.JsonElement; +import com.google.gson.stream.JsonReader; +import java.io.BufferedReader; +import java.io.IOException; +import javax.servlet.http.HttpServletRequest; +import org.skywalking.apm.collector.agentstream.jetty.handler.reader.TraceSegment; +import org.skywalking.apm.collector.agentstream.jetty.handler.reader.TraceSegmentJsonReader; +import org.skywalking.apm.collector.agentstream.worker.segment.SegmentParse; +import org.skywalking.apm.collector.server.jetty.ArgumentsParseException; +import org.skywalking.apm.collector.server.jetty.JettyHandler; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class TraceSegmentServletHandler extends JettyHandler { + + private final Logger logger = LoggerFactory.getLogger(TraceSegmentServletHandler.class); + + @Override public String pathSpec() { + return "/segments"; + } + + @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { + throw new UnsupportedOperationException(); + } + + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { + logger.debug("receive stream segment"); + try { + BufferedReader bufferedReader = req.getReader(); + read(bufferedReader); + } catch (IOException e) { + logger.error(e.getMessage(), e); + } + return null; + } + + private TraceSegmentJsonReader jsonReader = new TraceSegmentJsonReader(); + + private void read(BufferedReader bufferedReader) throws IOException { + JsonReader reader = new JsonReader(bufferedReader); + reader.beginArray(); + while (reader.hasNext()) { + SegmentParse segmentParse = new SegmentParse(); + TraceSegment traceSegment = jsonReader.read(reader); + segmentParse.parse(traceSegment.getGlobalTraceIds(), traceSegment.getTraceSegmentObject()); + } + reader.endArray(); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/KeyWithStringValueJsonReader.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/KeyWithStringValueJsonReader.java new file mode 100644 index 000000000..f0f1bb144 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/KeyWithStringValueJsonReader.java @@ -0,0 +1,36 @@ +package org.skywalking.apm.collector.agentstream.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; +import org.skywalking.apm.network.proto.KeyWithStringValue; + +/** + * @author pengys5 + */ +public class KeyWithStringValueJsonReader implements StreamJsonReader { + + private static final String KEY = "k"; + private static final String VALUE = "v"; + + @Override public KeyWithStringValue read(JsonReader reader) throws IOException { + KeyWithStringValue.Builder builder = KeyWithStringValue.newBuilder(); + + reader.beginObject(); + while (reader.hasNext()) { + switch (reader.nextName()) { + case KEY: + builder.setKey(reader.nextString()); + break; + case VALUE: + builder.setValue(reader.nextString()); + break; + default: + reader.skipValue(); + break; + } + } + reader.endObject(); + + return builder.build(); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/LogJsonReader.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/LogJsonReader.java new file mode 100644 index 000000000..f8d6526b2 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/LogJsonReader.java @@ -0,0 +1,37 @@ +package org.skywalking.apm.collector.agentstream.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; +import org.skywalking.apm.network.proto.LogMessage; + +/** + * @author pengys5 + */ +public class LogJsonReader implements StreamJsonReader { + + private KeyWithStringValueJsonReader keyWithStringValueJsonReader = new KeyWithStringValueJsonReader(); + + private static final String TI = "ti"; + private static final String LD = "ld"; + + @Override public LogMessage read(JsonReader reader) throws IOException { + LogMessage.Builder builder = LogMessage.newBuilder(); + + while (reader.hasNext()) { + switch (reader.nextName()) { + case TI: + builder.setTime(reader.nextLong()); + case LD: + reader.beginArray(); + while (reader.hasNext()) { + builder.addData(keyWithStringValueJsonReader.read(reader)); + } + reader.endArray(); + default: + reader.skipValue(); + } + } + + return builder.build(); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/ReferenceJsonReader.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/ReferenceJsonReader.java new file mode 100644 index 000000000..1398c4f53 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/ReferenceJsonReader.java @@ -0,0 +1,70 @@ +package org.skywalking.apm.collector.agentstream.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; +import org.skywalking.apm.network.proto.TraceSegmentReference; + +/** + * @author pengys5 + */ +public class ReferenceJsonReader implements StreamJsonReader { + + private UniqueIdJsonReader uniqueIdJsonReader = new UniqueIdJsonReader(); + + private static final String TS = "ts"; + private static final String AI = "ai"; + private static final String SI = "si"; + private static final String VI = "vi"; + private static final String VN = "vn"; + private static final String NI = "ni"; + private static final String NN = "nn"; + private static final String EI = "ei"; + private static final String EN = "en"; + private static final String RV = "rv"; + + @Override public TraceSegmentReference read(JsonReader reader) throws IOException { + TraceSegmentReference.Builder builder = TraceSegmentReference.newBuilder(); + + reader.beginObject(); + while (reader.hasNext()) { + switch (reader.nextName()) { + case TS: + builder.setParentTraceSegmentId(uniqueIdJsonReader.read(reader)); + break; + case AI: + builder.setParentApplicationInstanceId(reader.nextInt()); + break; + case SI: + builder.setParentSpanId(reader.nextInt()); + break; + case VI: + builder.setParentServiceId(reader.nextInt()); + break; + case VN: + builder.setParentServiceName(reader.nextString()); + break; + case NI: + builder.setNetworkAddressId(reader.nextInt()); + break; + case NN: + builder.setNetworkAddress(reader.nextString()); + break; + case EI: + builder.setEntryServiceId(reader.nextInt()); + break; + case EN: + builder.setEntryServiceName(reader.nextString()); + break; + case RV: + builder.setRefTypeValue(reader.nextInt()); + break; + default: + reader.skipValue(); + break; + } + } + reader.endObject(); + + return builder.build(); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/SegmentJsonReader.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/SegmentJsonReader.java new file mode 100644 index 000000000..c88a18f4b --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/SegmentJsonReader.java @@ -0,0 +1,69 @@ +package org.skywalking.apm.collector.agentstream.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; +import org.skywalking.apm.network.proto.TraceSegmentObject; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class SegmentJsonReader implements StreamJsonReader { + + private final Logger logger = LoggerFactory.getLogger(SegmentJsonReader.class); + + private UniqueIdJsonReader uniqueIdJsonReader = new UniqueIdJsonReader(); + private ReferenceJsonReader referenceJsonReader = new ReferenceJsonReader(); + private SpanJsonReader spanJsonReader = new SpanJsonReader(); + + private static final String TS = "ts"; + private static final String AI = "ai"; + private static final String II = "ii"; + private static final String RS = "rs"; + private static final String SS = "ss"; + + @Override public TraceSegmentObject read(JsonReader reader) throws IOException { + TraceSegmentObject.Builder builder = TraceSegmentObject.newBuilder(); + + reader.beginObject(); + while (reader.hasNext()) { + switch (reader.nextName()) { + case TS: + builder.setTraceSegmentId(uniqueIdJsonReader.read(reader)); + if (logger.isDebugEnabled()) { + StringBuilder segmentId = new StringBuilder(); + builder.getTraceSegmentId().getIdPartsList().forEach(idPart -> segmentId.append(idPart)); + logger.debug("segment id: {}", segmentId); + } + break; + case AI: + builder.setApplicationId(reader.nextInt()); + break; + case II: + builder.setApplicationInstanceId(reader.nextInt()); + break; + case RS: + reader.beginArray(); + while (reader.hasNext()) { + builder.addRefs(referenceJsonReader.read(reader)); + } + reader.endArray(); + break; + case SS: + reader.beginArray(); + while (reader.hasNext()) { + builder.addSpans(spanJsonReader.read(reader)); + } + reader.endArray(); + break; + default: + reader.skipValue(); + break; + } + } + reader.endObject(); + + return builder.build(); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/SpanJsonReader.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/SpanJsonReader.java new file mode 100644 index 000000000..074929525 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/SpanJsonReader.java @@ -0,0 +1,99 @@ +package org.skywalking.apm.collector.agentstream.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; +import org.skywalking.apm.network.proto.SpanObject; + +/** + * @author pengys5 + */ +public class SpanJsonReader implements StreamJsonReader { + + private KeyWithStringValueJsonReader keyWithStringValueJsonReader = new KeyWithStringValueJsonReader(); + private LogJsonReader logJsonReader = new LogJsonReader(); + + private static final String SI = "si"; + private static final String TV = "tv"; + private static final String LV = "lv"; + private static final String PS = "ps"; + private static final String ST = "st"; + private static final String ET = "et"; + private static final String CI = "ci"; + private static final String CN = "cn"; + private static final String OI = "oi"; + private static final String ON = "on"; + private static final String PI = "pi"; + private static final String PN = "pn"; + private static final String IE = "ie"; + private static final String TO = "to"; + private static final String LO = "lo"; + + @Override public SpanObject read(JsonReader reader) throws IOException { + SpanObject.Builder builder = SpanObject.newBuilder(); + + reader.beginObject(); + while (reader.hasNext()) { + switch (reader.nextName()) { + case SI: + builder.setSpanId(reader.nextInt()); + break; + case TV: + builder.setSpanTypeValue(reader.nextInt()); + break; + case LV: + builder.setSpanLayerValue(reader.nextInt()); + break; + case PS: + builder.setParentSpanId(reader.nextInt()); + break; + case ST: + builder.setStartTime(reader.nextLong()); + break; + case ET: + builder.setEndTime(reader.nextLong()); + break; + case CI: + builder.setComponentId(reader.nextInt()); + break; + case CN: + builder.setComponent(reader.nextString()); + break; + case OI: + builder.setOperationNameId(reader.nextInt()); + break; + case ON: + builder.setOperationName(reader.nextString()); + break; + case PI: + builder.setPeerId(reader.nextInt()); + break; + case PN: + builder.setPeer(reader.nextString()); + break; + case IE: + builder.setIsError(reader.nextBoolean()); + break; + case TO: + reader.beginArray(); + while (reader.hasNext()) { + builder.addTags(keyWithStringValueJsonReader.read(reader)); + } + reader.endArray(); + break; + case LO: + reader.beginArray(); + while (reader.hasNext()) { + builder.addLogs(logJsonReader.read(reader)); + } + reader.endArray(); + break; + default: + reader.skipValue(); + break; + } + } + reader.endObject(); + + return builder.build(); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/StreamJsonReader.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/StreamJsonReader.java new file mode 100644 index 000000000..95095a731 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/StreamJsonReader.java @@ -0,0 +1,11 @@ +package org.skywalking.apm.collector.agentstream.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; + +/** + * @author pengys5 + */ +public interface StreamJsonReader { + T read(JsonReader reader) throws IOException; +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/TraceSegment.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/TraceSegment.java new file mode 100644 index 000000000..fa56b0f04 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/TraceSegment.java @@ -0,0 +1,35 @@ +package org.skywalking.apm.collector.agentstream.jetty.handler.reader; + +import java.util.ArrayList; +import java.util.List; +import org.skywalking.apm.network.proto.TraceSegmentObject; +import org.skywalking.apm.network.proto.UniqueId; + +/** + * @author pengys5 + */ +public class TraceSegment { + + private List uniqueIds; + private TraceSegmentObject traceSegmentObject; + + public TraceSegment() { + uniqueIds = new ArrayList<>(); + } + + public List getGlobalTraceIds() { + return uniqueIds; + } + + public void addGlobalTraceId(UniqueId globalTraceId) { + uniqueIds.add(globalTraceId); + } + + public TraceSegmentObject getTraceSegmentObject() { + return traceSegmentObject; + } + + public void setTraceSegmentObject(TraceSegmentObject traceSegmentObject) { + this.traceSegmentObject = traceSegmentObject; + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/TraceSegmentJsonReader.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/TraceSegmentJsonReader.java new file mode 100644 index 000000000..26045e3c9 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/TraceSegmentJsonReader.java @@ -0,0 +1,54 @@ +package org.skywalking.apm.collector.agentstream.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class TraceSegmentJsonReader implements StreamJsonReader { + + private final Logger logger = LoggerFactory.getLogger(TraceSegmentJsonReader.class); + + private UniqueIdJsonReader uniqueIdJsonReader = new UniqueIdJsonReader(); + private SegmentJsonReader segmentJsonReader = new SegmentJsonReader(); + + private static final String GT = "gt"; + private static final String SG = "sg"; + + @Override public TraceSegment read(JsonReader reader) throws IOException { + TraceSegment traceSegment = new TraceSegment(); + + reader.beginObject(); + while (reader.hasNext()) { + switch (reader.nextName()) { + case GT: + reader.beginArray(); + while (reader.hasNext()) { + traceSegment.addGlobalTraceId(uniqueIdJsonReader.read(reader)); + } + reader.endArray(); + + if (logger.isDebugEnabled()) { + traceSegment.getGlobalTraceIds().forEach(uniqueId -> { + StringBuilder globalTraceId = new StringBuilder(); + uniqueId.getIdPartsList().forEach(idPart -> globalTraceId.append(idPart)); + logger.debug("global trace id: {}", globalTraceId.toString()); + }); + } + break; + case SG: + traceSegment.setTraceSegmentObject(segmentJsonReader.read(reader)); + break; + default: + reader.skipValue(); + break; + } + } + reader.endObject(); + + return traceSegment; + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/UniqueIdJsonReader.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/UniqueIdJsonReader.java new file mode 100644 index 000000000..8df8e1e0b --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/UniqueIdJsonReader.java @@ -0,0 +1,22 @@ +package org.skywalking.apm.collector.agentstream.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; +import org.skywalking.apm.network.proto.UniqueId; + +/** + * @author pengys5 + */ +public class UniqueIdJsonReader implements StreamJsonReader { + + @Override public UniqueId read(JsonReader reader) throws IOException { + UniqueId.Builder builder = UniqueId.newBuilder(); + + reader.beginArray(); + while (reader.hasNext()) { + builder.addIdParts(reader.nextLong()); + } + reader.endArray(); + return builder.build(); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/Const.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/Const.java index b0d8ec415..df48bd837 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/Const.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/Const.java @@ -8,6 +8,7 @@ public class Const { public static final String IDS_SPLIT = "\\.\\.-\\.\\."; public static final String PEERS_FRONT_SPLIT = "["; public static final String PEERS_BEHIND_SPLIT = "]"; + public static final int USER_ID = 1; public static final String USER_CODE = "User"; public static final String SEGMENT_SPAN_SPLIT = "S"; } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java new file mode 100644 index 000000000..54ac23020 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java @@ -0,0 +1,25 @@ +package org.skywalking.apm.collector.agentstream.worker.cache; + +import com.google.common.cache.Cache; +import com.google.common.cache.CacheBuilder; +import org.skywalking.apm.collector.agentstream.worker.register.instance.dao.IInstanceDAO; +import org.skywalking.apm.collector.storage.dao.DAOContainer; + +/** + * @author pengys5 + */ +public class InstanceCache { + + private static Cache CACHE = CacheBuilder.newBuilder().maximumSize(1000).build(); + + public static int get(int applicationInstanceId) { + try { + return CACHE.get(applicationInstanceId, () -> { + IInstanceDAO dao = (IInstanceDAO)DAOContainer.INSTANCE.get(IInstanceDAO.class.getName()); + return dao.getApplicationId(applicationInstanceId); + }); + } catch (Throwable e) { + return 0; + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentSpanListener.java index 4f4691802..6f83321d4 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/NodeComponentSpanListener.java @@ -40,7 +40,7 @@ public class NodeComponentSpanListener implements EntrySpanListener, ExitSpanLis peer = spanObject.getPeer(); } - String agg = componentName + Const.ID_SPLIT + peer; + String agg = peer + Const.ID_SPLIT + componentName; nodeComponents.add(agg); } @@ -62,7 +62,7 @@ public class NodeComponentSpanListener implements EntrySpanListener, ExitSpanLis } String peer = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId); - String agg = componentName + Const.ID_SPLIT + peer; + String agg = peer + Const.ID_SPLIT + componentName; nodeComponents.add(agg); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java index 2c5ebe600..38230c29b 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/NodeMappingSpanListener.java @@ -32,11 +32,11 @@ public class NodeMappingSpanListener implements RefsListener, FirstSpanListener String segmentId) { logger.debug("node mapping listener parse reference"); String peers = reference.getNetworkAddress(); - if (reference.getNetworkAddressId() == 0) { + if (reference.getNetworkAddressId() != 0) { peers = ExchangeMarkUtils.INSTANCE.buildMarkedID(reference.getNetworkAddressId()); } - String agg = applicationId + Const.ID_SPLIT + peers; + String agg = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId) + Const.ID_SPLIT + peers; nodeMappings.add(agg); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefSpanListener.java index a57e678b1..ba89e24be 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/NodeRefSpanListener.java @@ -3,6 +3,7 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.reference; import java.util.ArrayList; import java.util.List; import org.skywalking.apm.collector.agentstream.worker.Const; +import org.skywalking.apm.collector.agentstream.worker.cache.InstanceCache; import org.skywalking.apm.collector.agentstream.worker.noderef.reference.define.NodeRefDataDefine; import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener; @@ -27,27 +28,27 @@ public class NodeRefSpanListener implements EntrySpanListener, ExitSpanListener, private final Logger logger = LoggerFactory.getLogger(NodeRefSpanListener.class); - private List nodeExitReferences = new ArrayList<>(); + private List nodeReferences = new ArrayList<>(); private List nodeEntryReferences = new ArrayList<>(); private long timeBucket; private boolean hasReference = false; @Override public void parseExit(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { - String front = String.valueOf(applicationId); + String front = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId); String behind = spanObject.getPeer(); - if (spanObject.getPeerId() == 0) { + if (spanObject.getPeerId() != 0) { behind = ExchangeMarkUtils.INSTANCE.buildMarkedID(spanObject.getPeerId()); } String agg = front + Const.ID_SPLIT + behind; - nodeExitReferences.add(agg); + nodeReferences.add(agg); } @Override public void parseEntry(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { String behind = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId); - String front = Const.USER_CODE; + String front = ExchangeMarkUtils.INSTANCE.buildMarkedID(Const.USER_ID); String agg = front + Const.ID_SPLIT + behind; nodeEntryReferences.add(agg); } @@ -59,6 +60,14 @@ public class NodeRefSpanListener implements EntrySpanListener, ExitSpanListener, @Override public void parseRef(TraceSegmentReference reference, int applicationId, int applicationInstanceId, String segmentId) { + int parentApplicationId = InstanceCache.get(reference.getParentApplicationInstanceId()); + + String front = ExchangeMarkUtils.INSTANCE.buildMarkedID(parentApplicationId); + String behind = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId); + + String agg = front + Const.ID_SPLIT + behind; + nodeReferences.add(agg); + hasReference = true; } @@ -66,10 +75,10 @@ public class NodeRefSpanListener implements EntrySpanListener, ExitSpanListener, logger.debug("node reference listener build"); StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); if (!hasReference) { - nodeExitReferences.addAll(nodeEntryReferences); + nodeReferences.addAll(nodeEntryReferences); } - for (String agg : nodeExitReferences) { + for (String agg : nodeReferences) { NodeRefDataDefine.NodeReference nodeReference = new NodeRefDataDefine.NodeReference(); nodeReference.setId(timeBucket + Const.ID_SPLIT + agg); nodeReference.setAgg(agg); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java index f3251d032..aedfdeceb 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java @@ -3,6 +3,7 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.summary; import java.util.ArrayList; import java.util.List; import org.skywalking.apm.collector.agentstream.worker.Const; +import org.skywalking.apm.collector.agentstream.worker.cache.InstanceCache; import org.skywalking.apm.collector.agentstream.worker.noderef.summary.define.NodeRefSumDataDefine; import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener; @@ -29,14 +30,18 @@ public class NodeRefSumSpanListener implements EntrySpanListener, ExitSpanListen private List nodeExitReferences = new ArrayList<>(); private List nodeEntryReferences = new ArrayList<>(); + private List nodeReferences = new ArrayList<>(); private long timeBucket; private boolean hasReference = false; + private long startTime; + private long endTime; + private boolean isError; @Override public void parseExit(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { - String front = String.valueOf(applicationId); + String front = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId); String behind = spanObject.getPeer(); - if (spanObject.getPeerId() == 0) { + if (spanObject.getPeerId() != 0) { behind = ExchangeMarkUtils.INSTANCE.buildMarkedID(spanObject.getPeerId()); } @@ -47,8 +52,7 @@ public class NodeRefSumSpanListener implements EntrySpanListener, ExitSpanListen @Override public void parseEntry(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { String behind = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId); - String front = Const.USER_CODE; - + String front = ExchangeMarkUtils.INSTANCE.buildMarkedID(Const.USER_ID); String agg = front + Const.ID_SPLIT + behind; nodeEntryReferences.add(buildNodeRefSum(spanObject.getStartTime(), spanObject.getEndTime(), agg, spanObject.getIsError())); } @@ -77,11 +81,22 @@ public class NodeRefSumSpanListener implements EntrySpanListener, ExitSpanListen @Override public void parseFirst(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { timeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(spanObject.getStartTime()); + startTime = spanObject.getStartTime(); + endTime = spanObject.getEndTime(); + isError = spanObject.getIsError(); } @Override public void parseRef(TraceSegmentReference reference, int applicationId, int applicationInstanceId, String segmentId) { + int parentApplicationId = InstanceCache.get(reference.getParentApplicationInstanceId()); + + String front = ExchangeMarkUtils.INSTANCE.buildMarkedID(parentApplicationId); + String behind = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId); + + String agg = front + Const.ID_SPLIT + behind; + hasReference = true; + nodeReferences.add(agg); } @Override public void build() { @@ -89,6 +104,10 @@ public class NodeRefSumSpanListener implements EntrySpanListener, ExitSpanListen StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); if (!hasReference) { nodeExitReferences.addAll(nodeEntryReferences); + } else { + nodeReferences.forEach(agg -> { + nodeExitReferences.add(buildNodeRefSum(startTime, endTime, agg, isError)); + }); } for (NodeRefSumDataDefine.NodeReferenceSum referenceSum : nodeExitReferences) { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationRegisterSerialWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationRegisterSerialWorker.java index 791b3768e..ba2811a42 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationRegisterSerialWorker.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationRegisterSerialWorker.java @@ -1,5 +1,6 @@ package org.skywalking.apm.collector.agentstream.worker.register.application; +import org.skywalking.apm.collector.agentstream.worker.Const; import org.skywalking.apm.collector.agentstream.worker.register.IdAutoIncrement; import org.skywalking.apm.collector.agentstream.worker.register.application.dao.IApplicationDAO; import org.skywalking.apm.collector.storage.dao.DAOContainer; @@ -40,8 +41,11 @@ public class ApplicationRegisterSerialWorker extends AbstractLocalAsyncWorker { if (applicationId == 0) { int min = dao.getMinApplicationId(); if (min == 0) { - application.setApplicationId(1); - application.setId("1"); + ApplicationDataDefine.Application userApplication = new ApplicationDataDefine.Application(String.valueOf(Const.USER_ID), Const.USER_CODE, Const.USER_ID); + dao.save(userApplication); + + application.setApplicationId(-1); + application.setId("-1"); } else { int max = dao.getMaxApplicationId(); applicationId = IdAutoIncrement.INSTANCE.increment(min, max); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/IInstanceDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/IInstanceDAO.java index f0c311146..79d22d27e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/IInstanceDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/IInstanceDAO.java @@ -15,4 +15,6 @@ public interface IInstanceDAO { void save(InstanceDataDefine.Instance instance); void updateHeartbeatTime(int instanceId, long heartbeatTime); + + int getApplicationId(int applicationInstanceId); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceEsDAO.java index 237f50009..3b56987d9 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceEsDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceEsDAO.java @@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.register.instance.dao; import java.util.HashMap; import java.util.Map; +import org.elasticsearch.action.get.GetResponse; import org.elasticsearch.action.index.IndexResponse; import org.elasticsearch.action.search.SearchRequestBuilder; import org.elasticsearch.action.search.SearchResponse; @@ -40,8 +41,7 @@ public class InstanceEsDAO extends EsDAO implements IInstanceDAO { SearchResponse searchResponse = searchRequestBuilder.execute().actionGet(); if (searchResponse.getHits().totalHits > 0) { SearchHit searchHit = searchResponse.getHits().iterator().next(); - int instanceId = (int)searchHit.getSource().get(InstanceTable.COLUMN_INSTANCE_ID); - return instanceId; + return (int)searchHit.getSource().get(InstanceTable.COLUMN_INSTANCE_ID); } return 0; } @@ -57,7 +57,7 @@ public class InstanceEsDAO extends EsDAO implements IInstanceDAO { @Override public void save(InstanceDataDefine.Instance instance) { logger.debug("save instance register info, application id: {}, agentUUID: {}", instance.getApplicationId(), instance.getAgentUUID()); ElasticSearchClient client = getClient(); - Map source = new HashMap(); + Map source = new HashMap<>(); source.put(InstanceTable.COLUMN_INSTANCE_ID, instance.getInstanceId()); source.put(InstanceTable.COLUMN_APPLICATION_ID, instance.getApplicationId()); source.put(InstanceTable.COLUMN_AGENTUUID, instance.getAgentUUID()); @@ -75,10 +75,19 @@ public class InstanceEsDAO extends EsDAO implements IInstanceDAO { updateRequest.id(String.valueOf(instanceId)); updateRequest.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE); - Map source = new HashMap(); + Map source = new HashMap<>(); source.put(InstanceTable.COLUMN_HEARTBEAT_TIME, heartbeatTime); updateRequest.doc(source); client.update(updateRequest); } + + @Override public int getApplicationId(int applicationInstanceId) { + GetResponse response = getClient().prepareGet(InstanceTable.TABLE, String.valueOf(applicationInstanceId)).get(); + if (response.isExists()) { + return (int)response.getSource().get(InstanceTable.COLUMN_APPLICATION_ID); + } else { + return 0; + } + } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceH2DAO.java index 64070090e..d46aaaaec 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceH2DAO.java @@ -26,4 +26,8 @@ public class InstanceH2DAO extends H2DAO implements IInstanceDAO { @Override public void updateHeartbeatTime(int instanceId, long heartbeatTime) { } + + @Override public int getApplicationId(int applicationInstanceId) { + return 0; + } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/storage/PersistenceTimer.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/storage/PersistenceTimer.java index aeb38e243..6d17a4dfc 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/storage/PersistenceTimer.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/storage/PersistenceTimer.java @@ -30,19 +30,25 @@ public class PersistenceTimer implements Starter { } private void extractDataAndSave() { - List workers = PersistenceWorkerContainer.INSTANCE.getPersistenceWorkers(); - List batchAllCollection = new ArrayList<>(); - workers.forEach((PersistenceWorker worker) -> { - try { - worker.allocateJob(new FlushAndSwitch()); - List batchCollection = worker.buildBatchCollection(); - batchAllCollection.addAll(batchCollection); - } catch (WorkerException e) { - logger.error(e.getMessage(), e); - } - }); + try { + List workers = PersistenceWorkerContainer.INSTANCE.getPersistenceWorkers(); + List batchAllCollection = new ArrayList<>(); + workers.forEach((PersistenceWorker worker) -> { + try { + worker.allocateJob(new FlushAndSwitch()); + List batchCollection = worker.buildBatchCollection(); + batchAllCollection.addAll(batchCollection); + } catch (WorkerException e) { + logger.error(e.getMessage(), e); + } + }); - IBatchDAO dao = (IBatchDAO)DAOContainer.INSTANCE.get(IBatchDAO.class.getName()); - dao.batchPersistence(batchAllCollection); + IBatchDAO dao = (IBatchDAO)DAOContainer.INSTANCE.get(IBatchDAO.class.getName()); + dao.batchPersistence(batchAllCollection); + } catch (Throwable e) { + logger.error(e.getMessage(), e); + } finally { + logger.debug("persistence data save finish"); + } } } diff --git a/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/HttpClientTools.java b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/HttpClientTools.java new file mode 100644 index 000000000..ccaf9184c --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/HttpClientTools.java @@ -0,0 +1,77 @@ +package org.skywalking.apm.collector.agentstream; + +import java.io.IOException; +import java.net.URI; +import java.util.List; +import org.apache.http.Consts; +import org.apache.http.HttpEntity; +import org.apache.http.NameValuePair; +import org.apache.http.client.entity.UrlEncodedFormEntity; +import org.apache.http.client.methods.CloseableHttpResponse; +import org.apache.http.client.methods.HttpGet; +import org.apache.http.client.methods.HttpPost; +import org.apache.http.entity.StringEntity; +import org.apache.http.impl.client.CloseableHttpClient; +import org.apache.http.impl.client.HttpClients; +import org.apache.http.util.EntityUtils; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public enum HttpClientTools { + INSTANCE; + + private final Logger logger = LoggerFactory.getLogger(HttpClientTools.class); + + public String get(String url, List params) throws IOException { + CloseableHttpClient httpClient = HttpClients.createDefault(); + try { + HttpGet httpget = new HttpGet(url); + String paramStr = EntityUtils.toString(new UrlEncodedFormEntity(params)); + httpget.setURI(new URI(httpget.getURI().toString() + "?" + paramStr)); + logger.debug("executing get request {}", httpget.getURI()); + + try (CloseableHttpResponse response = httpClient.execute(httpget)) { + HttpEntity entity = response.getEntity(); + if (entity != null) { + return EntityUtils.toString(entity); + } + } + } catch (Exception e) { + logger.error(e.getMessage(), e); + } finally { + try { + httpClient.close(); + } catch (IOException e) { + logger.error(e.getMessage(), e); + } + } + return null; + } + + public String post(String url, String data) throws IOException { + CloseableHttpClient httpClient = HttpClients.createDefault(); + try { + HttpPost httppost = new HttpPost(url); + httppost.setEntity(new StringEntity(data, Consts.UTF_8)); + logger.debug("executing post request {}", httppost.getURI()); + try (CloseableHttpResponse response = httpClient.execute(httppost)) { + HttpEntity entity = response.getEntity(); + if (entity != null) { + return EntityUtils.toString(entity); + } + } + } catch (Exception e) { + logger.error(e.getMessage(), e); + } finally { + try { + httpClient.close(); + } catch (Exception e) { + logger.error(e.getMessage(), e); + } + } + return null; + } +} diff --git a/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/TraceSegmentJsonReaderTestCase.java b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/TraceSegmentJsonReaderTestCase.java new file mode 100644 index 000000000..5e42bec96 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/jetty/handler/reader/TraceSegmentJsonReaderTestCase.java @@ -0,0 +1,28 @@ +package org.skywalking.apm.collector.agentstream.jetty.handler.reader; + +import com.google.gson.JsonElement; +import com.google.gson.stream.JsonReader; +import java.io.IOException; +import java.io.StringReader; +import org.junit.Test; +import org.skywalking.apm.collector.agentstream.mock.JsonFileReader; + +/** + * @author pengys5 + */ +public class TraceSegmentJsonReaderTestCase { + + @Test + public void testRead() throws IOException { + TraceSegmentJsonReader reader = new TraceSegmentJsonReader(); + JsonElement jsonElement = JsonFileReader.INSTANCE.read("json/segment/normal/dubbox-consumer.json"); + System.out.println(jsonElement.toString()); + + JsonReader jsonReader = new JsonReader(new StringReader(jsonElement.toString())); + jsonReader.beginArray(); + while (jsonReader.hasNext()) { + reader.read(jsonReader); + } + jsonReader.endArray(); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/JsonFileReader.java b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/JsonFileReader.java new file mode 100644 index 000000000..b22c8edfc --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/JsonFileReader.java @@ -0,0 +1,24 @@ +package org.skywalking.apm.collector.agentstream.mock; + +import com.google.gson.JsonElement; +import com.google.gson.JsonParser; +import java.io.FileNotFoundException; +import java.io.FileReader; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public enum JsonFileReader { + INSTANCE; + + private final Logger logger = LoggerFactory.getLogger(JsonFileReader.class); + + public JsonElement read(String fileName) throws FileNotFoundException { + String path = this.getClass().getClassLoader().getResource(fileName).getFile(); + logger.debug("path: {}", path); + JsonParser jsonParser = new JsonParser(); + return jsonParser.parse(new FileReader(path)); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java new file mode 100644 index 000000000..80228cd16 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java @@ -0,0 +1,39 @@ +package org.skywalking.apm.collector.agentstream.mock; + +import com.google.gson.JsonElement; +import java.io.IOException; +import org.skywalking.apm.collector.agentstream.HttpClientTools; +import org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceDataDefine; +import org.skywalking.apm.collector.agentstream.worker.register.instance.dao.InstanceEsDAO; +import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; +import org.skywalking.apm.collector.core.CollectorException; + +/** + * @author pengys5 + */ +public class SegmentPost { + + // @Test + public void test() throws IOException, InterruptedException, CollectorException { + ElasticSearchClient client = new ElasticSearchClient("CollectorDBCluster", true, "127.0.0.1:9300"); + client.initialize(); + + InstanceEsDAO instanceEsDAO = new InstanceEsDAO(); + instanceEsDAO.setClient(client); + + InstanceDataDefine.Instance consumerInstance = new InstanceDataDefine.Instance("2", 2, "dubbox-consumer", 1501858094526L, 2); + instanceEsDAO.save(consumerInstance); + InstanceDataDefine.Instance providerInstance = new InstanceDataDefine.Instance("3", 3, "dubbox-provider", 1501858094526L, 3); + instanceEsDAO.save(providerInstance); + + JsonElement consumer = JsonFileReader.INSTANCE.read("json/segment/normal/dubbox-consumer.json"); + HttpClientTools.INSTANCE.post("http://localhost:12800/segments", consumer.toString()); + + Thread.sleep(5000); + + JsonElement provider = JsonFileReader.INSTANCE.read("json/segment/normal/dubbox-provider.json"); + HttpClientTools.INSTANCE.post("http://localhost:12800/segments", provider.toString()); + + Thread.sleep(5000); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/application-register.json b/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/application-register.json new file mode 100644 index 000000000..eb829c2cb --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/application-register.json @@ -0,0 +1,4 @@ +[ + "dubbox-consumer", + "dubbox-provider" +] \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/dubbox-consumer.json b/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/dubbox-consumer.json new file mode 100644 index 000000000..57c79c9c3 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/dubbox-consumer.json @@ -0,0 +1,71 @@ +[ + { + "gt": [ + [ + 230150, + 185809, + 24040000 + ] + ], + "sg": { + "ts": [ + 230150, + 185809, + 24040000 + ], + "ai": 2, + "ii": 2, + "rs": [], + "ss": [ + { + "si": 1, + "tv": 1, + "lv": 1, + "ps": 0, + "st": 1501858094526, + "et": 1501858097004, + "ci": 3, + "cn": "", + "oi": 0, + "on": "org.skywaking.apm.testcase.dubbo.services.GreetService.doBusiness()", + "pi": 0, + "pn": "172.25.0.4:20880", + "ie": false, + "to": [ + { + "k": "url", + "v": "rest://172.25.0.4:20880/org.skywaking.apm.testcase.dubbo.services.GreetService.doBusiness()" + } + ], + "lo": [] + }, + { + "si": 0, + "tv": 0, + "lv": 2, + "ps": -1, + "st": 1501858092409, + "et": 1501858097033, + "ci": 1, + "cn": "", + "oi": 0, + "on": "/dubbox-case/case/dubbox-rest", + "pi": 0, + "pn": "", + "ie": false, + "to": [ + { + "k": "url", + "v": "http://localhost:18080/dubbox-case/case/dubbox-rest" + }, + { + "k": "http.method", + "v": "GET" + } + ], + "lo": [] + } + ] + } + } +] \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/dubbox-provider.json b/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/dubbox-provider.json new file mode 100644 index 000000000..0e47af810 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/dubbox-provider.json @@ -0,0 +1,66 @@ +[ + { + "gt": [ + [ + 230150, + 185809, + 24040000 + ] + ], + "sg": { + "ts": [ + 137150, + 185809, + 48780000 + ], + "ai": 3, + "ii": 3, + "rs": [ + { + "ts": [ + 230150, + 185809, + 24040000 + ], + "ai": 2, + "si": 1, + "vi": 0, + "vn": "org.skywaking.apm.testcase.dubbo.services.GreetService.doBusiness()", + "ni": 0, + "nn": "172.25.0.4:20880", + "ei": 0, + "en": "/dubbox-case/case/dubbox-rest", + "rn": 0 + } + ], + "ss": [ + { + "si": 0, + "tv": 0, + "lv": 1, + "ps": -1, + "st": 1501858094883, + "et": 1501858096950, + "ci": 3, + "cn": "", + "oi": 0, + "on": "org.skywaking.apm.testcase.dubbo.services.GreetService.doBusiness()", + "pi": 0, + "pn": "", + "ie": false, + "to": [ + { + "k": "url", + "v": "rest://172.25.0.4:20880/org.skywaking.apm.testcase.dubbo.services.GreetService.doBusiness()" + }, + { + "k": "http.method", + "v": "GET" + } + ], + "lo": [] + } + ] + } + } +] \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/instance-register.json b/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/instance-register.json new file mode 100644 index 000000000..f29ab121d --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/test/resources/json/segment/normal/instance-register.json @@ -0,0 +1,5 @@ +{ + "ai": 0, + "au": "dubbox-consumer", + "rt": 1501858094526 +} \ No newline at end of file diff --git a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java index 1e4560af3..c2362ff5e 100644 --- a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java +++ b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java @@ -112,7 +112,6 @@ public class ElasticSearchClient implements Client { } public IndexRequestBuilder prepareIndex(String indexName, String id) { - client.prepareUpdate(); return client.prepareIndex(indexName, "type", id); } diff --git a/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/JettyHandler.java b/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/JettyHandler.java index 7921df426..22e21afba 100644 --- a/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/JettyHandler.java +++ b/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/JettyHandler.java @@ -13,6 +13,7 @@ import javax.servlet.http.HttpServlet; import javax.servlet.http.HttpServletRequest; import javax.servlet.http.HttpServletResponse; import org.skywalking.apm.collector.core.framework.Handler; +import org.skywalking.apm.collector.core.util.ObjectUtils; /** * @author pengys5 @@ -35,14 +36,13 @@ public abstract class JettyHandler extends HttpServlet implements Handler { @Override protected final void doPost(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException { try { - doPost(req); - reply(resp); + reply(resp, doPost(req)); } catch (ArgumentsParseException e) { replyError(resp, e.getMessage(), HttpServletResponse.SC_BAD_REQUEST); } } - protected abstract void doPost(HttpServletRequest req) throws ArgumentsParseException; + protected abstract JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException; @Override protected final void doHead(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException { @@ -136,17 +136,9 @@ public abstract class JettyHandler extends HttpServlet implements Handler { response.setStatus(HttpServletResponse.SC_OK); PrintWriter out = response.getWriter(); - out.print(resJson); - out.flush(); - out.close(); - } - - private void reply(HttpServletResponse response) throws IOException { - response.setContentType("text/json"); - response.setCharacterEncoding("utf-8"); - response.setStatus(HttpServletResponse.SC_OK); - - PrintWriter out = response.getWriter(); + if (ObjectUtils.isNotEmpty(resJson)) { + out.print(resJson); + } out.flush(); out.close(); } @@ -161,4 +153,4 @@ public abstract class JettyHandler extends HttpServlet implements Handler { out.flush(); out.close(); } -} +} \ No newline at end of file diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/BatchEsDAO.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/BatchEsDAO.java index 74c06e634..5a222145a 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/BatchEsDAO.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/BatchEsDAO.java @@ -4,6 +4,7 @@ import java.util.List; import org.elasticsearch.action.bulk.BulkRequestBuilder; import org.elasticsearch.action.bulk.BulkResponse; import org.elasticsearch.action.index.IndexRequestBuilder; +import org.elasticsearch.action.update.UpdateRequestBuilder; import org.skywalking.apm.collector.core.util.CollectionUtils; import org.skywalking.apm.collector.storage.dao.IBatchDAO; import org.slf4j.Logger; @@ -19,11 +20,16 @@ public class BatchEsDAO extends EsDAO implements IBatchDAO { @Override public void batchPersistence(List batchCollection) { BulkRequestBuilder bulkRequest = getClient().prepareBulk(); - logger.info("bulk data size: {}", batchCollection.size()); + logger.debug("bulk data size: {}", batchCollection.size()); if (CollectionUtils.isNotEmpty(batchCollection)) { for (int i = 0; i < batchCollection.size(); i++) { - IndexRequestBuilder builder = (IndexRequestBuilder)batchCollection.get(i); - bulkRequest.add(builder); + Object builder = batchCollection.get(i); + if (builder instanceof IndexRequestBuilder) { + bulkRequest.add((IndexRequestBuilder)builder); + } + if (builder instanceof UpdateRequestBuilder) { + bulkRequest.add((UpdateRequestBuilder)builder); + } } BulkResponse bulkResponse = bulkRequest.execute().actionGet(); diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/SegmentTopGetHandler.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/SegmentTopGetHandler.java index 195638203..003d37e55 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/SegmentTopGetHandler.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/SegmentTopGetHandler.java @@ -80,7 +80,7 @@ public class SegmentTopGetHandler extends JettyHandler { return service.loadTop(startTime, endTime, minCost, maxCost, operationName, globalTraceId, limit, from); } - @Override protected void doPost(HttpServletRequest req) throws ArgumentsParseException { + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { throw new UnsupportedOperationException(); } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/SpanGetHandler.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/SpanGetHandler.java index db1fa3a64..241e5bea5 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/SpanGetHandler.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/SpanGetHandler.java @@ -36,7 +36,7 @@ public class SpanGetHandler extends JettyHandler { return service.load(segmentId, spanId); } - @Override protected void doPost(HttpServletRequest req) throws ArgumentsParseException { + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { throw new UnsupportedOperationException(); } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/TraceDagGetHandler.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/TraceDagGetHandler.java index 494b004ae..a11540d83 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/TraceDagGetHandler.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/TraceDagGetHandler.java @@ -32,7 +32,7 @@ public class TraceDagGetHandler extends JettyHandler { return service.load(startTime, endTime, timeBucketType); } - @Override protected void doPost(HttpServletRequest req) throws ArgumentsParseException { + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { throw new UnsupportedOperationException(); } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/TraceStackGetHandler.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/TraceStackGetHandler.java index c24eb53af..b497e2580 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/TraceStackGetHandler.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/TraceStackGetHandler.java @@ -28,7 +28,7 @@ public class TraceStackGetHandler extends JettyHandler { return service.load(globalTraceId); } - @Override protected void doPost(HttpServletRequest req) throws ArgumentsParseException { + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { throw new UnsupportedOperationException(); } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/UIJettyServerHandler.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/UIJettyServerHandler.java index c53e5112f..78b49cd42 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/UIJettyServerHandler.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/handler/UIJettyServerHandler.java @@ -31,7 +31,7 @@ public class UIJettyServerHandler extends JettyHandler { return serverArray; } - @Override protected void doPost(HttpServletRequest req) throws ArgumentsParseException { + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { throw new UnsupportedOperationException(); } }