From 119d21d5252b466a269dbc793f8b6186d96dc8ce Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Sun, 6 Aug 2017 23:26:17 +0800 Subject: [PATCH 1/4] Service entry test success. --- .../noderef/summary/dao/NodeRefSumEsDAO.java | 14 +-- .../worker/segment/SegmentParse.java | 2 + .../entry/ServiceEntryAggregationWorker.java | 66 +++++++++++ .../entry/ServiceEntryPersistenceWorker.java | 71 ++++++++++++ .../entry/ServiceEntryRemoteWorker.java | 60 ++++++++++ .../entry/ServiceEntrySpanListener.java | 70 ++++++++++++ .../service/entry/dao/IServiceEntryDAO.java | 7 ++ .../service/entry/dao/ServiceEntryEsDAO.java | 50 +++++++++ .../service/entry/dao/ServiceEntryH2DAO.java | 9 ++ .../entry/define/ServiceEntryDataDefine.java | 106 ++++++++++++++++++ .../define/ServiceEntryEsTableDefine.java | 31 +++++ .../define/ServiceEntryH2TableDefine.java | 20 ++++ .../entry/define/ServiceEntryTable.java | 11 ++ .../worker/storage/PersistenceTimer.java | 1 + .../resources/META-INF/defines/es_dao.define | 3 +- .../resources/META-INF/defines/h2_dao.define | 3 +- .../defines/local_worker_provider.define | 3 + .../defines/remote_worker_provider.define | 4 +- .../resources/META-INF/defines/storage.define | 5 +- .../agentstream/mock/SegmentPost.java | 3 +- 20 files changed, 526 insertions(+), 13 deletions(-) create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntryAggregationWorker.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntryPersistenceWorker.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntryRemoteWorker.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/ServiceEntrySpanListener.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/IServiceEntryDAO.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/ServiceEntryEsDAO.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/ServiceEntryH2DAO.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryDataDefine.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryEsTableDefine.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryH2TableDefine.java create mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryTable.java 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(); From 4d7ad3377181d675547cc7cc37052719a3f3627c Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Mon, 7 Aug 2017 22:59:11 +0800 Subject: [PATCH 2/4] add entryApplicationInstanceId attribute that let collector known the entry service which application instance send. --- apm-network/src/main/proto/TraceSegmentService.proto | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/apm-network/src/main/proto/TraceSegmentService.proto b/apm-network/src/main/proto/TraceSegmentService.proto index a1e0aa134..ea4c2c994 100644 --- a/apm-network/src/main/proto/TraceSegmentService.proto +++ b/apm-network/src/main/proto/TraceSegmentService.proto @@ -35,10 +35,11 @@ message TraceSegmentReference { int32 parentApplicationInstanceId = 4; string networkAddress = 5; int32 networkAddressId = 6; - string entryServiceName = 7; - int32 entryServiceId = 8; - string parentServiceName = 9; - int32 parentServiceId = 10; + int32 entryApplicationInstanceId = 7; + string entryServiceName = 8; + int32 entryServiceId = 9; + string parentServiceName = 10; + int32 parentServiceId = 11; } message SpanObject { From eb72cb04bab47e65113b53c7d212d786081d3e28 Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Tue, 8 Aug 2017 09:41:25 +0800 Subject: [PATCH 3/4] Change the static json key name to full name in json reader. --- .../jetty/handler/reader/LogJsonReader.java | 8 +-- .../handler/reader/ReferenceJsonReader.java | 44 +++++++------- .../handler/reader/SegmentJsonReader.java | 20 +++---- .../jetty/handler/reader/SpanJsonReader.java | 60 +++++++++---------- .../reader/TraceSegmentJsonReader.java | 8 +-- 5 files changed, 70 insertions(+), 70 deletions(-) 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 index f8d6526b2..78b8a79ee 100644 --- 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 @@ -11,17 +11,17 @@ public class LogJsonReader implements StreamJsonReader { private KeyWithStringValueJsonReader keyWithStringValueJsonReader = new KeyWithStringValueJsonReader(); - private static final String TI = "ti"; - private static final String LD = "ld"; + private static final String TIME = "ti"; + private static final String LOG_DATA = "ld"; @Override public LogMessage read(JsonReader reader) throws IOException { LogMessage.Builder builder = LogMessage.newBuilder(); while (reader.hasNext()) { switch (reader.nextName()) { - case TI: + case TIME: builder.setTime(reader.nextLong()); - case LD: + case LOG_DATA: reader.beginArray(); while (reader.hasNext()) { builder.addData(keyWithStringValueJsonReader.read(reader)); 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 index bd979d6be..113edd199 100644 --- 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 @@ -11,17 +11,17 @@ public class ReferenceJsonReader implements StreamJsonReader { 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"; + private static final String TRACE_SEGMENT_ID = "ts"; + private static final String APPLICATION_ID = "ai"; + private static final String APPLICATION_INSTANCE_ID = "ii"; + private static final String TRACE_SEGMENT_REFERENCE = "rs"; + private static final String SPANS = "ss"; @Override public TraceSegmentObject read(JsonReader reader) throws IOException { TraceSegmentObject.Builder builder = TraceSegmentObject.newBuilder(); @@ -29,7 +29,7 @@ public class SegmentJsonReader implements StreamJsonReader { reader.beginObject(); while (reader.hasNext()) { switch (reader.nextName()) { - case TS: + case TRACE_SEGMENT_ID: builder.setTraceSegmentId(uniqueIdJsonReader.read(reader)); if (logger.isDebugEnabled()) { StringBuilder segmentId = new StringBuilder(); @@ -37,20 +37,20 @@ public class SegmentJsonReader implements StreamJsonReader { logger.debug("segment id: {}", segmentId); } break; - case AI: + case APPLICATION_ID: builder.setApplicationId(reader.nextInt()); break; - case II: + case APPLICATION_INSTANCE_ID: builder.setApplicationInstanceId(reader.nextInt()); break; - case RS: + case TRACE_SEGMENT_REFERENCE: reader.beginArray(); while (reader.hasNext()) { builder.addRefs(referenceJsonReader.read(reader)); } reader.endArray(); break; - case SS: + case SPANS: reader.beginArray(); while (reader.hasNext()) { builder.addSpans(spanJsonReader.read(reader)); 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 index 074929525..8de743512 100644 --- 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 @@ -12,21 +12,21 @@ 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"; + private static final String SPAN_ID = "si"; + private static final String SPAN_TYPE_VALUE = "tv"; + private static final String SPAN_LAYER_VALUE = "lv"; + private static final String PARENT_SPAN_ID = "ps"; + private static final String START_TIME = "st"; + private static final String END_TIME = "et"; + private static final String COMPONENT_ID = "ci"; + private static final String COMPONENT_NAME = "cn"; + private static final String OPERATION_NAME_ID = "oi"; + private static final String OPERATION_NAME = "on"; + private static final String PEER_ID = "pi"; + private static final String PEER = "pn"; + private static final String IS_ERROR = "ie"; + private static final String TAGS = "to"; + private static final String LOGS = "lo"; @Override public SpanObject read(JsonReader reader) throws IOException { SpanObject.Builder builder = SpanObject.newBuilder(); @@ -34,53 +34,53 @@ public class SpanJsonReader implements StreamJsonReader { reader.beginObject(); while (reader.hasNext()) { switch (reader.nextName()) { - case SI: + case SPAN_ID: builder.setSpanId(reader.nextInt()); break; - case TV: + case SPAN_TYPE_VALUE: builder.setSpanTypeValue(reader.nextInt()); break; - case LV: + case SPAN_LAYER_VALUE: builder.setSpanLayerValue(reader.nextInt()); break; - case PS: + case PARENT_SPAN_ID: builder.setParentSpanId(reader.nextInt()); break; - case ST: + case START_TIME: builder.setStartTime(reader.nextLong()); break; - case ET: + case END_TIME: builder.setEndTime(reader.nextLong()); break; - case CI: + case COMPONENT_ID: builder.setComponentId(reader.nextInt()); break; - case CN: + case COMPONENT_NAME: builder.setComponent(reader.nextString()); break; - case OI: + case OPERATION_NAME_ID: builder.setOperationNameId(reader.nextInt()); break; - case ON: + case OPERATION_NAME: builder.setOperationName(reader.nextString()); break; - case PI: + case PEER_ID: builder.setPeerId(reader.nextInt()); break; - case PN: + case PEER: builder.setPeer(reader.nextString()); break; - case IE: + case IS_ERROR: builder.setIsError(reader.nextBoolean()); break; - case TO: + case TAGS: reader.beginArray(); while (reader.hasNext()) { builder.addTags(keyWithStringValueJsonReader.read(reader)); } reader.endArray(); break; - case LO: + case LOGS: reader.beginArray(); while (reader.hasNext()) { builder.addLogs(logJsonReader.read(reader)); 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 index 26045e3c9..d00883b53 100644 --- 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 @@ -15,8 +15,8 @@ public class TraceSegmentJsonReader implements StreamJsonReader { private UniqueIdJsonReader uniqueIdJsonReader = new UniqueIdJsonReader(); private SegmentJsonReader segmentJsonReader = new SegmentJsonReader(); - private static final String GT = "gt"; - private static final String SG = "sg"; + private static final String GLOBAL_TRACE_IDS = "gt"; + private static final String SEGMENT = "sg"; @Override public TraceSegment read(JsonReader reader) throws IOException { TraceSegment traceSegment = new TraceSegment(); @@ -24,7 +24,7 @@ public class TraceSegmentJsonReader implements StreamJsonReader { reader.beginObject(); while (reader.hasNext()) { switch (reader.nextName()) { - case GT: + case GLOBAL_TRACE_IDS: reader.beginArray(); while (reader.hasNext()) { traceSegment.addGlobalTraceId(uniqueIdJsonReader.read(reader)); @@ -39,7 +39,7 @@ public class TraceSegmentJsonReader implements StreamJsonReader { }); } break; - case SG: + case SEGMENT: traceSegment.setTraceSegmentObject(segmentJsonReader.read(reader)); break; default: From a7c11a1baf9a2bd0e25295d12cedac872bd2a0d0 Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Tue, 8 Aug 2017 10:03:34 +0800 Subject: [PATCH 4/4] delete the define id in every data define. #336 #338 --- .../worker/global/define/GlobalTraceDataDefine.java | 4 ---- .../component/define/NodeComponentDataDefine.java | 4 ---- .../node/mapping/define/NodeMappingDataDefine.java | 4 ---- .../noderef/reference/define/NodeRefDataDefine.java | 4 ---- .../summary/define/NodeRefSumDataDefine.java | 6 ------ .../register/application/ApplicationDataDefine.java | 8 +------- .../register/application/ApplicationTable.java | 4 +++- .../register/instance/InstanceDataDefine.java | 10 ++-------- .../register/instance/InstanceEsTableDefine.java | 2 +- .../register/instance/InstanceH2TableDefine.java | 2 +- .../worker/register/instance/InstanceTable.java | 6 ++++-- .../worker/register/instance/dao/InstanceEsDAO.java | 4 ++-- .../register/servicename/ServiceNameDataDefine.java | 8 +------- .../register/servicename/ServiceNameTable.java | 4 +++- .../segment/cost/define/SegmentCostDataDefine.java | 4 ---- .../segment/origin/define/SegmentDataDefine.java | 6 ------ .../entry/define/ServiceEntryDataDefine.java | 4 ---- .../reference/define/ServiceRefDataDefine.java | 6 ------ .../core/cluster/ClusterDefinitionFile.java | 13 ------------- .../apm/collector/stream/worker/impl/data/Data.java | 8 +------- .../stream/worker/impl/data/DataDefine.java | 4 +--- 21 files changed, 20 insertions(+), 95 deletions(-) delete mode 100644 apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterDefinitionFile.java diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceDataDefine.java index 16ba94c1e..9d851d5e5 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceDataDefine.java @@ -14,10 +14,6 @@ import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation */ public class GlobalTraceDataDefine extends DataDefine { - @Override public int defineId() { - return 403; - } - @Override protected int initialCapacity() { return 4; } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentDataDefine.java index 0f08484ce..e617a9679 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentDataDefine.java @@ -14,10 +14,6 @@ import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation */ public class NodeComponentDataDefine extends DataDefine { - @Override public int defineId() { - return 101; - } - @Override protected int initialCapacity() { return 3; } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingDataDefine.java index e227b291e..881a0f163 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingDataDefine.java @@ -14,10 +14,6 @@ import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation */ public class NodeMappingDataDefine extends DataDefine { - @Override public int defineId() { - return 102; - } - @Override protected int initialCapacity() { return 3; } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefDataDefine.java index 079062861..7ec0fb367 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/reference/define/NodeRefDataDefine.java @@ -14,10 +14,6 @@ import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation */ public class NodeRefDataDefine extends DataDefine { - @Override public int defineId() { - return 201; - } - @Override protected int initialCapacity() { return 3; } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumDataDefine.java index 39146ea63..192c6630e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumDataDefine.java @@ -14,12 +14,6 @@ import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation */ public class NodeRefSumDataDefine extends DataDefine { - public static final int DEFINE_ID = 202; - - @Override public int defineId() { - return DEFINE_ID; - } - @Override protected int initialCapacity() { return 9; } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationDataDefine.java index cb964a3a4..3aebfe46a 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationDataDefine.java @@ -12,18 +12,12 @@ import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation */ public class ApplicationDataDefine extends DataDefine { - public static final int DEFINE_ID = 101; - - @Override public int defineId() { - return DEFINE_ID; - } - @Override protected int initialCapacity() { return 3; } @Override protected void attributeDefine() { - addAttribute(0, new Attribute("id", AttributeType.STRING, new NonOperation())); + addAttribute(0, new Attribute(ApplicationTable.COLUMN_ID, AttributeType.STRING, new NonOperation())); addAttribute(1, new Attribute(ApplicationTable.COLUMN_APPLICATION_CODE, AttributeType.STRING, new CoverOperation())); addAttribute(2, new Attribute(ApplicationTable.COLUMN_APPLICATION_ID, AttributeType.INTEGER, new CoverOperation())); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationTable.java index 797b8d80a..6b067c2b1 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationTable.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationTable.java @@ -1,9 +1,11 @@ package org.skywalking.apm.collector.agentstream.worker.register.application; +import org.skywalking.apm.collector.agentstream.worker.CommonTable; + /** * @author pengys5 */ -public class ApplicationTable { +public class ApplicationTable extends CommonTable { public static final String TABLE = "application"; public static final String COLUMN_APPLICATION_CODE = "application_code"; 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/register/instance/InstanceDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceDataDefine.java index e14afa7b2..cdb986013 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceDataDefine.java @@ -12,20 +12,14 @@ import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation */ public class InstanceDataDefine extends DataDefine { - public static final int DEFINE_ID = 102; - - @Override public int defineId() { - return DEFINE_ID; - } - @Override protected int initialCapacity() { return 6; } @Override protected void attributeDefine() { - addAttribute(0, new Attribute("id", AttributeType.STRING, new NonOperation())); + addAttribute(0, new Attribute(InstanceTable.COLUMN_ID, AttributeType.STRING, new NonOperation())); addAttribute(1, new Attribute(InstanceTable.COLUMN_APPLICATION_ID, AttributeType.INTEGER, new CoverOperation())); - addAttribute(2, new Attribute(InstanceTable.COLUMN_AGENTUUID, AttributeType.STRING, new CoverOperation())); + addAttribute(2, new Attribute(InstanceTable.COLUMN_AGENT_UUID, AttributeType.STRING, new CoverOperation())); addAttribute(3, new Attribute(InstanceTable.COLUMN_REGISTER_TIME, AttributeType.LONG, new CoverOperation())); addAttribute(4, new Attribute(InstanceTable.COLUMN_INSTANCE_ID, AttributeType.INTEGER, new CoverOperation())); addAttribute(5, new Attribute(InstanceTable.COLUMN_HEARTBEAT_TIME, AttributeType.LONG, new CoverOperation())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceEsTableDefine.java index dffcfa9fb..55417fbaf 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceEsTableDefine.java @@ -26,7 +26,7 @@ public class InstanceEsTableDefine extends ElasticSearchTableDefine { @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(InstanceTable.COLUMN_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); - addColumn(new ElasticSearchColumnDefine(InstanceTable.COLUMN_AGENTUUID, ElasticSearchColumnDefine.Type.Keyword.name())); + addColumn(new ElasticSearchColumnDefine(InstanceTable.COLUMN_AGENT_UUID, ElasticSearchColumnDefine.Type.Keyword.name())); addColumn(new ElasticSearchColumnDefine(InstanceTable.COLUMN_REGISTER_TIME, ElasticSearchColumnDefine.Type.Long.name())); addColumn(new ElasticSearchColumnDefine(InstanceTable.COLUMN_INSTANCE_ID, ElasticSearchColumnDefine.Type.Integer.name())); addColumn(new ElasticSearchColumnDefine(InstanceTable.COLUMN_HEARTBEAT_TIME, ElasticSearchColumnDefine.Type.Long.name())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceH2TableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceH2TableDefine.java index 58ca3e626..17afcea8e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceH2TableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceH2TableDefine.java @@ -14,7 +14,7 @@ public class InstanceH2TableDefine extends H2TableDefine { @Override public void initialize() { addColumn(new H2ColumnDefine(InstanceTable.COLUMN_APPLICATION_ID, H2ColumnDefine.Type.Int.name())); - addColumn(new H2ColumnDefine(InstanceTable.COLUMN_AGENTUUID, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(InstanceTable.COLUMN_AGENT_UUID, H2ColumnDefine.Type.Varchar.name())); addColumn(new H2ColumnDefine(InstanceTable.COLUMN_REGISTER_TIME, H2ColumnDefine.Type.Bigint.name())); addColumn(new H2ColumnDefine(InstanceTable.COLUMN_INSTANCE_ID, H2ColumnDefine.Type.Int.name())); addColumn(new H2ColumnDefine(InstanceTable.COLUMN_HEARTBEAT_TIME, H2ColumnDefine.Type.Bigint.name())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceTable.java index 5a26adb6e..fa59f8220 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceTable.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceTable.java @@ -1,12 +1,14 @@ package org.skywalking.apm.collector.agentstream.worker.register.instance; +import org.skywalking.apm.collector.agentstream.worker.CommonTable; + /** * @author pengys5 */ -public class InstanceTable { +public class InstanceTable extends CommonTable { public static final String TABLE = "instance"; public static final String COLUMN_APPLICATION_ID = "application_id"; - public static final String COLUMN_AGENTUUID = "agent_uuid"; + public static final String COLUMN_AGENT_UUID = "agent_uuid"; public static final String COLUMN_REGISTER_TIME = "register_time"; public static final String COLUMN_INSTANCE_ID = "instance_id"; public static final String COLUMN_HEARTBEAT_TIME = "heartbeatTime"; 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 3b56987d9..bf872bca5 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 @@ -34,7 +34,7 @@ public class InstanceEsDAO extends EsDAO implements IInstanceDAO { searchRequestBuilder.setSearchType(SearchType.QUERY_THEN_FETCH); BoolQueryBuilder builder = QueryBuilders.boolQuery(); builder.must().add(QueryBuilders.termQuery(InstanceTable.COLUMN_APPLICATION_ID, applicationId)); - builder.must().add(QueryBuilders.termQuery(InstanceTable.COLUMN_AGENTUUID, agentUUID)); + builder.must().add(QueryBuilders.termQuery(InstanceTable.COLUMN_AGENT_UUID, agentUUID)); searchRequestBuilder.setQuery(builder); searchRequestBuilder.setSize(1); @@ -60,7 +60,7 @@ public class InstanceEsDAO extends EsDAO implements IInstanceDAO { 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()); + source.put(InstanceTable.COLUMN_AGENT_UUID, instance.getAgentUUID()); source.put(InstanceTable.COLUMN_REGISTER_TIME, instance.getRegisterTime()); IndexResponse response = client.prepareIndex(InstanceTable.TABLE, instance.getId()).setSource(source).setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE).get(); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameDataDefine.java index d30a15240..42463c27e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameDataDefine.java @@ -12,18 +12,12 @@ import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation */ public class ServiceNameDataDefine extends DataDefine { - public static final int DEFINE_ID = 103; - - @Override public int defineId() { - return DEFINE_ID; - } - @Override protected int initialCapacity() { return 4; } @Override protected void attributeDefine() { - addAttribute(0, new Attribute("id", AttributeType.STRING, new NonOperation())); + addAttribute(0, new Attribute(ServiceNameTable.COLUMN_ID, AttributeType.STRING, new NonOperation())); addAttribute(1, new Attribute(ServiceNameTable.COLUMN_SERVICE_NAME, AttributeType.STRING, new CoverOperation())); addAttribute(2, new Attribute(ServiceNameTable.COLUMN_APPLICATION_ID, AttributeType.INTEGER, new CoverOperation())); addAttribute(3, new Attribute(ServiceNameTable.COLUMN_SERVICE_ID, AttributeType.INTEGER, new CoverOperation())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameTable.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameTable.java index f15fd6075..f6c81e491 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameTable.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameTable.java @@ -1,9 +1,11 @@ package org.skywalking.apm.collector.agentstream.worker.register.servicename; +import org.skywalking.apm.collector.agentstream.worker.CommonTable; + /** * @author pengys5 */ -public class ServiceNameTable { +public class ServiceNameTable extends CommonTable { public static final String TABLE = "service_name"; public static final String COLUMN_SERVICE_NAME = "service_name"; 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/segment/cost/define/SegmentCostDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostDataDefine.java index 51f57fd43..fd9f85c03 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostDataDefine.java @@ -14,10 +14,6 @@ import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation */ public class SegmentCostDataDefine extends DataDefine { - @Override public int defineId() { - return 402; - } - @Override protected int initialCapacity() { return 8; } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/define/SegmentDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/define/SegmentDataDefine.java index 6edf0322f..7e2f75caa 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/define/SegmentDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/define/SegmentDataDefine.java @@ -15,12 +15,6 @@ import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation */ public class SegmentDataDefine extends DataDefine { - public static final int DEFINE_ID = 401; - - @Override public int defineId() { - return DEFINE_ID; - } - @Override protected int initialCapacity() { return 2; } 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 index ed044eb5d..ea3a37d75 100644 --- 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 @@ -14,10 +14,6 @@ import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation */ public class ServiceEntryDataDefine extends DataDefine { - @Override public int defineId() { - return 501; - } - @Override protected int initialCapacity() { return 4; } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefDataDefine.java index 44248a920..55e4b723d 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/reference/define/ServiceRefDataDefine.java @@ -14,12 +14,6 @@ import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation */ public class ServiceRefDataDefine extends DataDefine { - public static final int DEFINE_ID = 501; - - @Override public int defineId() { - return DEFINE_ID; - } - @Override protected int initialCapacity() { return 4; } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterDefinitionFile.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterDefinitionFile.java deleted file mode 100644 index aa1bb7b54..000000000 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterDefinitionFile.java +++ /dev/null @@ -1,13 +0,0 @@ -package org.skywalking.apm.collector.core.cluster; - -import org.skywalking.apm.collector.core.framework.DefinitionFile; - -/** - * @author pengys5 - */ -public class ClusterDefinitionFile extends DefinitionFile { - - @Override protected String fileName() { - return "cluster-configuration.define"; - } -} diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Data.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Data.java index 5d5e25af8..5a8832ba3 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Data.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/Data.java @@ -8,7 +8,6 @@ import org.skywalking.apm.collector.stream.worker.selector.AbstractHashMessage; * @author pengys5 */ public class Data extends AbstractHashMessage { - private int defineId; private final int stringCapacity; private final int longCapacity; private final int floatCapacity; @@ -22,10 +21,9 @@ public class Data extends AbstractHashMessage { private Boolean[] dataBooleans; private byte[][] dataBytes; - public Data(String id, int defineId, int stringCapacity, int longCapacity, int floatCapacity, int integerCapacity, + public Data(String id, int stringCapacity, int longCapacity, int floatCapacity, int integerCapacity, int booleanCapacity, int byteCapacity) { super(id); - this.defineId = defineId; this.dataStrings = new String[stringCapacity]; this.dataLongs = new Long[longCapacity]; this.dataFloats = new Float[floatCapacity]; @@ -92,10 +90,6 @@ public class Data extends AbstractHashMessage { return dataStrings[0]; } - public int getDefineId() { - return defineId; - } - public RemoteData serialize() { RemoteData.Builder builder = RemoteData.newBuilder(); builder.setIntegerCapacity(integerCapacity); diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefine.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefine.java index 254e742d6..ddbd9c2a3 100644 --- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefine.java +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/DataDefine.java @@ -42,14 +42,12 @@ public abstract class DataDefine { attributes[position] = attribute; } - public abstract int defineId(); - protected abstract int initialCapacity(); protected abstract void attributeDefine(); public final Data build(String id) { - return new Data(id, defineId(), stringCapacity, longCapacity, floatCapacity, integerCapacity, booleanCapacity, byteCapacity); + return new Data(id, stringCapacity, longCapacity, floatCapacity, integerCapacity, booleanCapacity, byteCapacity); } public void mergeData(Data newData, Data oldData) {