diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumEsDAO.java index c117e70a9..1810412d3 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumEsDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumEsDAO.java @@ -22,13 +22,13 @@ public class NodeRefSumEsDAO extends EsDAO implements INodeRefSumDAO, IPersisten if (getResponse.isExists()) { Data data = dataDefine.build(id); Map source = getResponse.getSource(); - data.setDataLong(0, (Long)source.get(NodeRefSumTable.COLUMN_ONE_SECOND_LESS)); - data.setDataLong(1, (Long)source.get(NodeRefSumTable.COLUMN_THREE_SECOND_LESS)); - data.setDataLong(2, (Long)source.get(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS)); - data.setDataLong(3, (Long)source.get(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER)); - data.setDataLong(4, (Long)source.get(NodeRefSumTable.COLUMN_ERROR)); - data.setDataLong(5, (Long)source.get(NodeRefSumTable.COLUMN_SUMMARY)); - data.setDataLong(6, (Long)source.get(NodeRefSumTable.COLUMN_TIME_BUCKET)); + data.setDataLong(0, ((Number)source.get(NodeRefSumTable.COLUMN_ONE_SECOND_LESS)).longValue()); + data.setDataLong(1, ((Number)source.get(NodeRefSumTable.COLUMN_THREE_SECOND_LESS)).longValue()); + data.setDataLong(2, ((Number)source.get(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS)).longValue()); + data.setDataLong(3, ((Number)source.get(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER)).longValue()); + data.setDataLong(4, ((Number)source.get(NodeRefSumTable.COLUMN_ERROR)).longValue()); + data.setDataLong(5, ((Number)source.get(NodeRefSumTable.COLUMN_SUMMARY)).longValue()); + data.setDataLong(6, ((Number)source.get(NodeRefSumTable.COLUMN_TIME_BUCKET)).longValue()); data.setDataString(1, (String)source.get(NodeRefSumTable.COLUMN_AGG)); return data; } else { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SegmentParse.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SegmentParse.java index 6558a2615..0b4d436b9 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SegmentParse.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/SegmentParse.java @@ -10,6 +10,7 @@ import org.skywalking.apm.collector.agentstream.worker.noderef.summary.NodeRefSu import org.skywalking.apm.collector.agentstream.worker.segment.cost.SegmentCostSpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.origin.SegmentPersistenceWorker; import org.skywalking.apm.collector.agentstream.worker.segment.origin.define.SegmentDataDefine; +import org.skywalking.apm.collector.agentstream.worker.service.entry.ServiceEntrySpanListener; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.core.util.CollectionUtils; import org.skywalking.apm.collector.stream.StreamModuleContext; @@ -41,6 +42,7 @@ public class SegmentParse { spanListeners.add(new NodeRefSumSpanListener()); spanListeners.add(new SegmentCostSpanListener()); spanListeners.add(new GlobalTraceSpanListener()); + spanListeners.add(new ServiceEntrySpanListener()); } public void parse(List traceIds, TraceSegmentObject segmentObject) { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntryAggregationWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntryAggregationWorker.java new file mode 100644 index 000000000..c74fce3f8 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntryAggregationWorker.java @@ -0,0 +1,66 @@ +package org.skywalking.apm.collector.agentstream.worker.service.entry; + +import org.skywalking.apm.collector.agentstream.worker.service.entry.define.ServiceEntryDataDefine; +import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorkerProvider; +import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext; +import org.skywalking.apm.collector.stream.worker.ProviderNotFoundException; +import org.skywalking.apm.collector.stream.worker.Role; +import org.skywalking.apm.collector.stream.worker.WorkerNotFoundException; +import org.skywalking.apm.collector.stream.worker.WorkerRefs; +import org.skywalking.apm.collector.stream.worker.impl.AggregationWorker; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; +import org.skywalking.apm.collector.stream.worker.selector.HashCodeSelector; +import org.skywalking.apm.collector.stream.worker.selector.WorkerSelector; + +/** + * @author pengys5 + */ +public class ServiceEntryAggregationWorker extends AggregationWorker { + + public ServiceEntryAggregationWorker(Role role, ClusterWorkerContext clusterContext) { + super(role, clusterContext); + } + + @Override public void preStart() throws ProviderNotFoundException { + super.preStart(); + } + + @Override protected WorkerRefs nextWorkRef(String id) throws WorkerNotFoundException { + return getClusterContext().lookup(ServiceEntryRemoteWorker.WorkerRole.INSTANCE); + } + + public static class Factory extends AbstractLocalAsyncWorkerProvider { + @Override + public Role role() { + return WorkerRole.INSTANCE; + } + + @Override + public ServiceEntryAggregationWorker workerInstance(ClusterWorkerContext clusterContext) { + return new ServiceEntryAggregationWorker(role(), clusterContext); + } + + @Override + public int queueSize() { + return 1024; + } + } + + public enum WorkerRole implements Role { + INSTANCE; + + @Override + public String roleName() { + return ServiceEntryAggregationWorker.class.getSimpleName(); + } + + @Override + public WorkerSelector workerSelector() { + return new HashCodeSelector(); + } + + @Override public DataDefine dataDefine() { + return new ServiceEntryDataDefine(); + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntryPersistenceWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntryPersistenceWorker.java new file mode 100644 index 000000000..a5caec87a --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntryPersistenceWorker.java @@ -0,0 +1,71 @@ +package org.skywalking.apm.collector.agentstream.worker.service.entry; + +import org.skywalking.apm.collector.agentstream.worker.service.entry.dao.IServiceEntryDAO; +import org.skywalking.apm.collector.agentstream.worker.service.entry.define.ServiceEntryDataDefine; +import org.skywalking.apm.collector.storage.dao.DAOContainer; +import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorkerProvider; +import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext; +import org.skywalking.apm.collector.stream.worker.ProviderNotFoundException; +import org.skywalking.apm.collector.stream.worker.Role; +import org.skywalking.apm.collector.stream.worker.impl.PersistenceWorker; +import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; +import org.skywalking.apm.collector.stream.worker.selector.HashCodeSelector; +import org.skywalking.apm.collector.stream.worker.selector.WorkerSelector; + +/** + * @author pengys5 + */ +public class ServiceEntryPersistenceWorker extends PersistenceWorker { + + public ServiceEntryPersistenceWorker(Role role, ClusterWorkerContext clusterContext) { + super(role, clusterContext); + } + + @Override public void preStart() throws ProviderNotFoundException { + super.preStart(); + } + + @Override protected boolean needMergeDBData() { + return true; + } + + @Override protected IPersistenceDAO persistenceDAO() { + return (IPersistenceDAO)DAOContainer.INSTANCE.get(IServiceEntryDAO.class.getName()); + } + + public static class Factory extends AbstractLocalAsyncWorkerProvider { + @Override + public Role role() { + return WorkerRole.INSTANCE; + } + + @Override + public ServiceEntryPersistenceWorker workerInstance(ClusterWorkerContext clusterContext) { + return new ServiceEntryPersistenceWorker(role(), clusterContext); + } + + @Override + public int queueSize() { + return 1024; + } + } + + public enum WorkerRole implements Role { + INSTANCE; + + @Override + public String roleName() { + return ServiceEntryPersistenceWorker.class.getSimpleName(); + } + + @Override + public WorkerSelector workerSelector() { + return new HashCodeSelector(); + } + + @Override public DataDefine dataDefine() { + return new ServiceEntryDataDefine(); + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntryRemoteWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntryRemoteWorker.java new file mode 100644 index 000000000..3ca0a8195 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntryRemoteWorker.java @@ -0,0 +1,60 @@ +package org.skywalking.apm.collector.agentstream.worker.service.entry; + +import org.skywalking.apm.collector.agentstream.worker.service.entry.define.ServiceEntryDataDefine; +import org.skywalking.apm.collector.stream.worker.AbstractRemoteWorker; +import org.skywalking.apm.collector.stream.worker.AbstractRemoteWorkerProvider; +import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext; +import org.skywalking.apm.collector.stream.worker.ProviderNotFoundException; +import org.skywalking.apm.collector.stream.worker.Role; +import org.skywalking.apm.collector.stream.worker.WorkerException; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; +import org.skywalking.apm.collector.stream.worker.selector.HashCodeSelector; +import org.skywalking.apm.collector.stream.worker.selector.WorkerSelector; + +/** + * @author pengys5 + */ +public class ServiceEntryRemoteWorker extends AbstractRemoteWorker { + + protected ServiceEntryRemoteWorker(Role role, ClusterWorkerContext clusterContext) { + super(role, clusterContext); + } + + @Override public void preStart() throws ProviderNotFoundException { + + } + + @Override protected void onWork(Object message) throws WorkerException { + getClusterContext().lookup(ServiceEntryPersistenceWorker.WorkerRole.INSTANCE).tell(message); + } + + public static class Factory extends AbstractRemoteWorkerProvider { + @Override + public Role role() { + return WorkerRole.INSTANCE; + } + + @Override + public ServiceEntryRemoteWorker workerInstance(ClusterWorkerContext clusterContext) { + return new ServiceEntryRemoteWorker(role(), clusterContext); + } + } + + public enum WorkerRole implements Role { + INSTANCE; + + @Override + public String roleName() { + return ServiceEntryRemoteWorker.class.getSimpleName(); + } + + @Override + public WorkerSelector workerSelector() { + return new HashCodeSelector(); + } + + @Override public DataDefine dataDefine() { + return new ServiceEntryDataDefine(); + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntrySpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntrySpanListener.java new file mode 100644 index 000000000..aaeba441b --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntrySpanListener.java @@ -0,0 +1,70 @@ +package org.skywalking.apm.collector.agentstream.worker.service.entry; + +import org.skywalking.apm.collector.agentstream.worker.Const; +import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener; +import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener; +import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener; +import org.skywalking.apm.collector.agentstream.worker.service.entry.define.ServiceEntryDataDefine; +import org.skywalking.apm.collector.agentstream.worker.util.ExchangeMarkUtils; +import org.skywalking.apm.collector.agentstream.worker.util.TimeBucketUtils; +import org.skywalking.apm.collector.core.framework.CollectorContextHelper; +import org.skywalking.apm.collector.stream.StreamModuleContext; +import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; +import org.skywalking.apm.collector.stream.worker.WorkerInvokeException; +import org.skywalking.apm.collector.stream.worker.WorkerNotFoundException; +import org.skywalking.apm.network.proto.SpanObject; +import org.skywalking.apm.network.proto.TraceSegmentReference; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class ServiceEntrySpanListener implements RefsListener, FirstSpanListener, EntrySpanListener { + + private final Logger logger = LoggerFactory.getLogger(ServiceEntrySpanListener.class); + + private long timeBucket; + private boolean hasReference = false; + private String agg; + private int applicationId; + + @Override + public void parseEntry(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { + String entryServiceName = spanObject.getOperationName(); + if (spanObject.getOperationNameId() != 0) { + entryServiceName = ExchangeMarkUtils.INSTANCE.buildMarkedID(spanObject.getOperationNameId()); + } + this.agg = ExchangeMarkUtils.INSTANCE.buildMarkedID(applicationId) + Const.ID_SPLIT + entryServiceName; + this.applicationId = applicationId; + } + + @Override public void parseRef(TraceSegmentReference reference, int applicationId, int applicationInstanceId, + String segmentId) { + hasReference = true; + } + + @Override + public void parseFirst(SpanObject spanObject, int applicationId, int applicationInstanceId, String segmentId) { + timeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(spanObject.getStartTime()); + } + + @Override public void build() { + logger.debug("entry service listener build"); + StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); + if (!hasReference) { + ServiceEntryDataDefine.ServiceEntry serviceEntry = new ServiceEntryDataDefine.ServiceEntry(); + serviceEntry.setId(timeBucket + Const.ID_SPLIT + agg); + serviceEntry.setApplicationId(applicationId); + serviceEntry.setAgg(agg); + serviceEntry.setTimeBucket(timeBucket); + + try { + logger.debug("send to service entry aggregation worker, id: {}", serviceEntry.getId()); + context.getClusterWorkerContext().lookup(ServiceEntryAggregationWorker.WorkerRole.INSTANCE).tell(serviceEntry.toData()); + } catch (WorkerInvokeException | WorkerNotFoundException e) { + logger.error(e.getMessage(), e); + } + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/IServiceEntryDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/IServiceEntryDAO.java new file mode 100644 index 000000000..fcfe8f572 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/IServiceEntryDAO.java @@ -0,0 +1,7 @@ +package org.skywalking.apm.collector.agentstream.worker.service.entry.dao; + +/** + * @author pengys5 + */ +public interface IServiceEntryDAO { +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/ServiceEntryEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/ServiceEntryEsDAO.java new file mode 100644 index 000000000..b35edd115 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/ServiceEntryEsDAO.java @@ -0,0 +1,50 @@ +package org.skywalking.apm.collector.agentstream.worker.service.entry.dao; + +import java.util.HashMap; +import java.util.Map; +import org.elasticsearch.action.get.GetResponse; +import org.elasticsearch.action.index.IndexRequestBuilder; +import org.elasticsearch.action.update.UpdateRequestBuilder; +import org.skywalking.apm.collector.agentstream.worker.service.entry.define.ServiceEntryTable; +import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; +import org.skywalking.apm.collector.stream.worker.impl.data.Data; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; + +/** + * @author pengys5 + */ +public class ServiceEntryEsDAO extends EsDAO implements IServiceEntryDAO, IPersistenceDAO { + + @Override public Data get(String id, DataDefine dataDefine) { + GetResponse getResponse = getClient().prepareGet(ServiceEntryTable.TABLE, id).get(); + if (getResponse.isExists()) { + Data data = dataDefine.build(id); + Map source = getResponse.getSource(); + data.setDataInteger(0, (Integer)source.get(ServiceEntryTable.COLUMN_APPLICATION_ID)); + data.setDataString(1, (String)source.get(ServiceEntryTable.COLUMN_AGG)); + data.setDataLong(0, (Long)source.get(ServiceEntryTable.COLUMN_TIME_BUCKET)); + return data; + } else { + return null; + } + } + + @Override public IndexRequestBuilder prepareBatchInsert(Data data) { + Map source = new HashMap<>(); + source.put(ServiceEntryTable.COLUMN_APPLICATION_ID, data.getDataInteger(0)); + source.put(ServiceEntryTable.COLUMN_AGG, data.getDataString(1)); + source.put(ServiceEntryTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + + return getClient().prepareIndex(ServiceEntryTable.TABLE, data.getDataString(0)).setSource(source); + } + + @Override public UpdateRequestBuilder prepareBatchUpdate(Data data) { + Map source = new HashMap<>(); + source.put(ServiceEntryTable.COLUMN_APPLICATION_ID, data.getDataInteger(0)); + source.put(ServiceEntryTable.COLUMN_AGG, data.getDataString(1)); + source.put(ServiceEntryTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + + return getClient().prepareUpdate(ServiceEntryTable.TABLE, data.getDataString(0)).setDoc(source); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/ServiceEntryH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/ServiceEntryH2DAO.java new file mode 100644 index 000000000..fb2cd96c7 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/ServiceEntryH2DAO.java @@ -0,0 +1,9 @@ +package org.skywalking.apm.collector.agentstream.worker.service.entry.dao; + +import org.skywalking.apm.collector.storage.h2.dao.H2DAO; + +/** + * @author pengys5 + */ +public class ServiceEntryH2DAO extends H2DAO implements IServiceEntryDAO { +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryDataDefine.java new file mode 100644 index 000000000..ed044eb5d --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryDataDefine.java @@ -0,0 +1,106 @@ +package org.skywalking.apm.collector.agentstream.worker.service.entry.define; + +import org.skywalking.apm.collector.remote.grpc.proto.RemoteData; +import org.skywalking.apm.collector.stream.worker.impl.data.Attribute; +import org.skywalking.apm.collector.stream.worker.impl.data.AttributeType; +import org.skywalking.apm.collector.stream.worker.impl.data.Data; +import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine; +import org.skywalking.apm.collector.stream.worker.impl.data.Transform; +import org.skywalking.apm.collector.stream.worker.impl.data.operate.CoverOperation; +import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation; + +/** + * @author pengys5 + */ +public class ServiceEntryDataDefine extends DataDefine { + + @Override public int defineId() { + return 501; + } + + @Override protected int initialCapacity() { + return 4; + } + + @Override protected void attributeDefine() { + addAttribute(0, new Attribute(ServiceEntryTable.COLUMN_ID, AttributeType.STRING, new NonOperation())); + addAttribute(1, new Attribute(ServiceEntryTable.COLUMN_APPLICATION_ID, AttributeType.INTEGER, new NonOperation())); + addAttribute(2, new Attribute(ServiceEntryTable.COLUMN_AGG, AttributeType.STRING, new CoverOperation())); + addAttribute(3, new Attribute(ServiceEntryTable.COLUMN_TIME_BUCKET, AttributeType.LONG, new CoverOperation())); + } + + @Override public Object deserialize(RemoteData remoteData) { + return null; + } + + @Override public RemoteData serialize(Object object) { + return null; + } + + public static class ServiceEntry implements Transform { + private String id; + private int applicationId; + private String agg; + private long timeBucket; + + ServiceEntry(String id, int applicationId, String agg, long timeBucket) { + this.id = id; + this.applicationId = applicationId; + this.agg = agg; + this.timeBucket = timeBucket; + } + + public ServiceEntry() { + } + + @Override public Data toData() { + ServiceEntryDataDefine define = new ServiceEntryDataDefine(); + Data data = define.build(id); + data.setDataString(0, this.id); + data.setDataInteger(0, this.applicationId); + data.setDataString(1, this.agg); + data.setDataLong(0, this.timeBucket); + return data; + } + + @Override public ServiceEntry toSelf(Data data) { + this.id = data.getDataString(0); + this.applicationId = data.getDataInteger(0); + this.agg = data.getDataString(1); + this.timeBucket = data.getDataLong(0); + return this; + } + + public String getId() { + return id; + } + + public String getAgg() { + return agg; + } + + public long getTimeBucket() { + return timeBucket; + } + + public void setId(String id) { + this.id = id; + } + + public void setAgg(String agg) { + this.agg = agg; + } + + public void setTimeBucket(long timeBucket) { + this.timeBucket = timeBucket; + } + + public int getApplicationId() { + return applicationId; + } + + public void setApplicationId(int applicationId) { + this.applicationId = applicationId; + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryEsTableDefine.java new file mode 100644 index 000000000..54191227f --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryEsTableDefine.java @@ -0,0 +1,31 @@ +package org.skywalking.apm.collector.agentstream.worker.service.entry.define; + +import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; +import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; + +/** + * @author pengys5 + */ +public class ServiceEntryEsTableDefine extends ElasticSearchTableDefine { + + public ServiceEntryEsTableDefine() { + super(ServiceEntryTable.TABLE); + } + + @Override public int refreshInterval() { + return 2; + } + + @Override public int numberOfShards() { + return 2; + } + + @Override public int numberOfReplicas() { + return 0; + } + + @Override public void initialize() { + addColumn(new ElasticSearchColumnDefine(ServiceEntryTable.COLUMN_AGG, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(ServiceEntryTable.COLUMN_TIME_BUCKET, ElasticSearchColumnDefine.Type.Long.name())); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryH2TableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryH2TableDefine.java new file mode 100644 index 000000000..e38d6e726 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryH2TableDefine.java @@ -0,0 +1,20 @@ +package org.skywalking.apm.collector.agentstream.worker.service.entry.define; + +import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; +import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; + +/** + * @author pengys5 + */ +public class ServiceEntryH2TableDefine extends H2TableDefine { + + public ServiceEntryH2TableDefine() { + super(ServiceEntryTable.TABLE); + } + + @Override public void initialize() { + addColumn(new H2ColumnDefine(ServiceEntryTable.COLUMN_ID, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(ServiceEntryTable.COLUMN_AGG, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(ServiceEntryTable.COLUMN_TIME_BUCKET, H2ColumnDefine.Type.Bigint.name())); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryTable.java new file mode 100644 index 000000000..5ec9957f8 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryTable.java @@ -0,0 +1,11 @@ +package org.skywalking.apm.collector.agentstream.worker.service.entry.define; + +import org.skywalking.apm.collector.agentstream.worker.CommonTable; + +/** + * @author pengys5 + */ +public class ServiceEntryTable extends CommonTable { + public static final String TABLE = "service_entry"; + public static final String COLUMN_APPLICATION_ID = "application_id"; +} 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 6d17a4dfc..92825552d 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 @@ -34,6 +34,7 @@ public class PersistenceTimer implements Starter { List workers = PersistenceWorkerContainer.INSTANCE.getPersistenceWorkers(); List batchAllCollection = new ArrayList<>(); workers.forEach((PersistenceWorker worker) -> { + logger.debug("extract {} worker data and save", worker.getRole().roleName()); try { worker.allocateJob(new FlushAndSwitch()); List batchCollection = worker.buildBatchCollection(); diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/es_dao.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/es_dao.define index dd2e4e545..0a9b11a86 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/es_dao.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/es_dao.define @@ -7,4 +7,5 @@ org.skywalking.apm.collector.agentstream.worker.noderef.reference.dao.NodeRefere org.skywalking.apm.collector.agentstream.worker.segment.origin.dao.SegmentEsDAO org.skywalking.apm.collector.agentstream.worker.noderef.summary.dao.NodeRefSumEsDAO org.skywalking.apm.collector.agentstream.worker.segment.cost.dao.SegmentCostEsDAO -org.skywalking.apm.collector.agentstream.worker.global.dao.GlobalTraceEsDAO \ No newline at end of file +org.skywalking.apm.collector.agentstream.worker.global.dao.GlobalTraceEsDAO +org.skywalking.apm.collector.agentstream.worker.service.entry.dao.ServiceEntryEsDAO \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/h2_dao.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/h2_dao.define index da2352e37..c02664268 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/h2_dao.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/h2_dao.define @@ -7,4 +7,5 @@ org.skywalking.apm.collector.agentstream.worker.noderef.reference.dao.NodeRefere org.skywalking.apm.collector.agentstream.worker.segment.origin.dao.SegmentH2DAO org.skywalking.apm.collector.agentstream.worker.noderef.summary.dao.NodeRefSumH2DAO org.skywalking.apm.collector.agentstream.worker.segment.cost.dao.SegmentCostH2DAO -org.skywalking.apm.collector.agentstream.worker.global.dao.GlobalTraceH2DAO \ No newline at end of file +org.skywalking.apm.collector.agentstream.worker.global.dao.GlobalTraceH2DAO +org.skywalking.apm.collector.agentstream.worker.service.entry.dao.ServiceEntryH2DAO \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_worker_provider.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_worker_provider.define index fb133fbe0..34e6bdd71 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_worker_provider.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_worker_provider.define @@ -10,6 +10,9 @@ org.skywalking.apm.collector.agentstream.worker.noderef.reference.NodeRefPersist org.skywalking.apm.collector.agentstream.worker.noderef.summary.NodeRefSumAggregationWorker$Factory org.skywalking.apm.collector.agentstream.worker.noderef.summary.NodeRefSumPersistenceWorker$Factory +org.skywalking.apm.collector.agentstream.worker.service.entry.ServiceEntryAggregationWorker$Factory +org.skywalking.apm.collector.agentstream.worker.service.entry.ServiceEntryPersistenceWorker$Factory + org.skywalking.apm.collector.agentstream.worker.segment.origin.SegmentPersistenceWorker$Factory org.skywalking.apm.collector.agentstream.worker.segment.cost.SegmentCostPersistenceWorker$Factory org.skywalking.apm.collector.agentstream.worker.global.GlobalTracePersistenceWorker$Factory diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/remote_worker_provider.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/remote_worker_provider.define index b82aef6bd..db8d6bb14 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/remote_worker_provider.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/remote_worker_provider.define @@ -5,4 +5,6 @@ org.skywalking.apm.collector.agentstream.worker.register.servicename.ServiceName org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentRemoteWorker$Factory org.skywalking.apm.collector.agentstream.worker.node.mapping.NodeMappingRemoteWorker$Factory org.skywalking.apm.collector.agentstream.worker.noderef.reference.NodeRefRemoteWorker$Factory -org.skywalking.apm.collector.agentstream.worker.noderef.summary.NodeRefSumRemoteWorker$Factory \ No newline at end of file +org.skywalking.apm.collector.agentstream.worker.noderef.summary.NodeRefSumRemoteWorker$Factory + +org.skywalking.apm.collector.agentstream.worker.service.entry.ServiceEntryRemoteWorker$Factory \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define index d046f4613..b5d27b7b1 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define @@ -26,4 +26,7 @@ org.skywalking.apm.collector.agentstream.worker.segment.cost.define.SegmentCostE org.skywalking.apm.collector.agentstream.worker.segment.cost.define.SegmentCostH2TableDefine org.skywalking.apm.collector.agentstream.worker.global.define.GlobalTraceEsTableDefine -org.skywalking.apm.collector.agentstream.worker.global.define.GlobalTraceH2TableDefine \ No newline at end of file +org.skywalking.apm.collector.agentstream.worker.global.define.GlobalTraceH2TableDefine + +org.skywalking.apm.collector.agentstream.worker.service.entry.define.ServiceEntryEsTableDefine +org.skywalking.apm.collector.agentstream.worker.service.entry.define.ServiceEntryH2TableDefine \ No newline at end of file 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 index 80228cd16..0a16508ab 100644 --- 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 @@ -13,8 +13,7 @@ import org.skywalking.apm.collector.core.CollectorException; */ public class SegmentPost { - // @Test - public void test() throws IOException, InterruptedException, CollectorException { + public static void main(String[] args) throws IOException, InterruptedException, CollectorException { ElasticSearchClient client = new ElasticSearchClient("CollectorDBCluster", true, "127.0.0.1:9300"); client.initialize();