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(); } }