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.
This commit is contained in:
pengys5 2017-08-06 21:54:55 +08:00
parent 3492b0d86e
commit 911de3fe4a
49 changed files with 1183 additions and 99 deletions

View File

@ -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;
}

View File

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

View File

@ -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() {
}
}

View File

@ -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<Handler> handlerList() {
List<Handler> handlers = new LinkedList<>();
handlers.add(new ApplicationRegisterServletHandler());
handlers.add(new InstanceDiscoveryServletHandler());
return handlers;
}
}

View File

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

View File

@ -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;
}
}

View File

@ -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;
}
}

View File

@ -1 +1,2 @@
org.skywalking.apm.collector.agentregister.grpc.AgentRegisterGRPCModuleDefine
org.skywalking.apm.collector.agentregister.grpc.AgentRegisterGRPCModuleDefine
org.skywalking.apm.collector.agentregister.jetty.AgentRegisterJettyModuleDefine

View File

@ -25,13 +25,11 @@ public class AgentStreamGRPCServerHandler extends JettyHandler {
ClusterModuleRegistrationReader reader = ((ClusterModuleContext)CollectorContextHelper.INSTANCE.getContext(ClusterModuleGroupDefine.GROUP_NAME)).getReader();
List<String> 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();
}
}

View File

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

View File

@ -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<Handler> handlerList() {
return null;
List<Handler> handlers = new LinkedList<>();
handlers.add(new TraceSegmentServletHandler());
return handlers;
}
}

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

@ -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> {
T read(JsonReader reader) throws IOException;
}

View File

@ -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<UniqueId> uniqueIds;
private TraceSegmentObject traceSegmentObject;
public TraceSegment() {
uniqueIds = new ArrayList<>();
}
public List<UniqueId> getGlobalTraceIds() {
return uniqueIds;
}
public void addGlobalTraceId(UniqueId globalTraceId) {
uniqueIds.add(globalTraceId);
}
public TraceSegmentObject getTraceSegmentObject() {
return traceSegmentObject;
}
public void setTraceSegmentObject(TraceSegmentObject traceSegmentObject) {
this.traceSegmentObject = traceSegmentObject;
}
}

View File

@ -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<TraceSegment> {
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;
}
}

View File

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

View File

@ -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";
}

View File

@ -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<Integer, Integer> 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;
}
}
}

View File

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

View File

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

View File

@ -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<String> nodeExitReferences = new ArrayList<>();
private List<String> nodeReferences = new ArrayList<>();
private List<String> 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);

View File

@ -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<NodeRefSumDataDefine.NodeReferenceSum> nodeExitReferences = new ArrayList<>();
private List<NodeRefSumDataDefine.NodeReferenceSum> nodeEntryReferences = new ArrayList<>();
private List<String> 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) {

View File

@ -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);

View File

@ -15,4 +15,6 @@ public interface IInstanceDAO {
void save(InstanceDataDefine.Instance instance);
void updateHeartbeatTime(int instanceId, long heartbeatTime);
int getApplicationId(int applicationInstanceId);
}

View File

@ -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<String, Object> source = new HashMap();
Map<String, Object> 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<String, Object> source = new HashMap();
Map<String, Object> 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;
}
}
}

View File

@ -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;
}
}

View File

@ -30,19 +30,25 @@ public class PersistenceTimer implements Starter {
}
private void extractDataAndSave() {
List<PersistenceWorker> 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<PersistenceWorker> 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");
}
}
}

View File

@ -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<NameValuePair> 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;
}
}

View File

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

View File

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

View File

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

View File

@ -0,0 +1,4 @@
[
"dubbox-consumer",
"dubbox-provider"
]

View File

@ -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": []
}
]
}
}
]

View File

@ -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": []
}
]
}
}
]

View File

@ -0,0 +1,5 @@
{
"ai": 0,
"au": "dubbox-consumer",
"rt": 1501858094526
}

View File

@ -112,7 +112,6 @@ public class ElasticSearchClient implements Client {
}
public IndexRequestBuilder prepareIndex(String indexName, String id) {
client.prepareUpdate();
return client.prepareIndex(indexName, "type", id);
}

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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