diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumAggregationWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumAggregationWorker.java new file mode 100644 index 000000000..4e6d4c8d8 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumAggregationWorker.java @@ -0,0 +1,66 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.summary; + +import org.skywalking.apm.collector.agentstream.worker.noderef.summary.define.NodeRefSumDataDefine; +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 NodeRefSumAggregationWorker extends AggregationWorker { + + public NodeRefSumAggregationWorker(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(NodeRefSumRemoteWorker.WorkerRole.INSTANCE); + } + + public static class Factory extends AbstractLocalAsyncWorkerProvider { + @Override + public Role role() { + return WorkerRole.INSTANCE; + } + + @Override + public NodeRefSumAggregationWorker workerInstance(ClusterWorkerContext clusterContext) { + return new NodeRefSumAggregationWorker(role(), clusterContext); + } + + @Override + public int queueSize() { + return 1024; + } + } + + public enum WorkerRole implements Role { + INSTANCE; + + @Override + public String roleName() { + return NodeRefSumAggregationWorker.class.getSimpleName(); + } + + @Override + public WorkerSelector workerSelector() { + return new HashCodeSelector(); + } + + @Override public DataDefine dataDefine() { + return new NodeRefSumDataDefine(); + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumPersistenceWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumPersistenceWorker.java new file mode 100644 index 000000000..96dead7e0 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumPersistenceWorker.java @@ -0,0 +1,70 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.summary; + +import java.util.List; +import java.util.Map; +import org.skywalking.apm.collector.agentstream.worker.noderef.summary.dao.INodeRefSumDAO; +import org.skywalking.apm.collector.agentstream.worker.noderef.summary.define.NodeRefSumDataDefine; +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.data.Data; +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 NodeRefSumPersistenceWorker extends PersistenceWorker { + + public NodeRefSumPersistenceWorker(Role role, ClusterWorkerContext clusterContext) { + super(role, clusterContext); + } + + @Override public void preStart() throws ProviderNotFoundException { + super.preStart(); + } + + @Override protected List prepareBatch(Map dataMap) { + INodeRefSumDAO dao = (INodeRefSumDAO)DAOContainer.INSTANCE.get(INodeRefSumDAO.class.getName()); + return dao.prepareBatch(dataMap); + } + + public static class Factory extends AbstractLocalAsyncWorkerProvider { + @Override + public Role role() { + return WorkerRole.INSTANCE; + } + + @Override + public NodeRefSumPersistenceWorker workerInstance(ClusterWorkerContext clusterContext) { + return new NodeRefSumPersistenceWorker(role(), clusterContext); + } + + @Override + public int queueSize() { + return 1024; + } + } + + public enum WorkerRole implements Role { + INSTANCE; + + @Override + public String roleName() { + return NodeRefSumPersistenceWorker.class.getSimpleName(); + } + + @Override + public WorkerSelector workerSelector() { + return new HashCodeSelector(); + } + + @Override public DataDefine dataDefine() { + return new NodeRefSumDataDefine(); + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumRemoteWorker.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumRemoteWorker.java new file mode 100644 index 000000000..5936faa68 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumRemoteWorker.java @@ -0,0 +1,60 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.summary; + +import org.skywalking.apm.collector.agentstream.worker.noderef.summary.define.NodeRefSumDataDefine; +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 NodeRefSumRemoteWorker extends AbstractRemoteWorker { + + protected NodeRefSumRemoteWorker(Role role, ClusterWorkerContext clusterContext) { + super(role, clusterContext); + } + + @Override public void preStart() throws ProviderNotFoundException { + + } + + @Override protected void onWork(Object message) throws WorkerException { + getClusterContext().lookup(NodeRefSumPersistenceWorker.WorkerRole.INSTANCE).tell(message); + } + + public static class Factory extends AbstractRemoteWorkerProvider { + @Override + public Role role() { + return WorkerRole.INSTANCE; + } + + @Override + public NodeRefSumRemoteWorker workerInstance(ClusterWorkerContext clusterContext) { + return new NodeRefSumRemoteWorker(role(), clusterContext); + } + } + + public enum WorkerRole implements Role { + INSTANCE; + + @Override + public String roleName() { + return NodeRefSumRemoteWorker.class.getSimpleName(); + } + + @Override + public WorkerSelector workerSelector() { + return new HashCodeSelector(); + } + + @Override public DataDefine dataDefine() { + return new NodeRefSumDataDefine(); + } + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java new file mode 100644 index 000000000..40a138dc8 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/NodeRefSumSpanListener.java @@ -0,0 +1,101 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.summary; + +import java.util.ArrayList; +import java.util.List; +import org.skywalking.apm.collector.agentstream.worker.Const; +import org.skywalking.apm.collector.agentstream.worker.noderef.summary.define.NodeRefSumDataDefine; +import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener; +import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener; +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.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 NodeRefSumSpanListener implements EntrySpanListener, ExitSpanListener, FirstSpanListener, RefsListener { + + private final Logger logger = LoggerFactory.getLogger(NodeRefSumSpanListener.class); + + private List nodeExitReferences = new ArrayList<>(); + private List nodeEntryReferences = new ArrayList<>(); + private long timeBucket; + private boolean hasReference = false; + + @Override public void parseExit(SpanObject spanObject, int applicationId, int applicationInstanceId) { + String front = String.valueOf(applicationId); + String behind = String.valueOf(spanObject.getPeerId()); + if (spanObject.getPeerId() == 0) { + behind = spanObject.getPeer(); + } + + String agg = front + Const.ID_SPLIT + behind; + nodeExitReferences.add(buildNodeRefSum(spanObject.getStartTime(), spanObject.getEndTime(), agg, spanObject.getIsError())); + } + + @Override public void parseEntry(SpanObject spanObject, int applicationId, int applicationInstanceId) { + String behind = String.valueOf(applicationId); + String front = Const.USER_CODE; + + String agg = front + Const.ID_SPLIT + behind; + nodeEntryReferences.add(buildNodeRefSum(spanObject.getStartTime(), spanObject.getEndTime(), agg, spanObject.getIsError())); + } + + private NodeRefSumDataDefine.NodeReferenceSum buildNodeRefSum(long startTime, long endTime, String agg, + boolean isError) { + NodeRefSumDataDefine.NodeReferenceSum referenceSum = new NodeRefSumDataDefine.NodeReferenceSum(); + referenceSum.setAgg(agg); + + long cost = endTime - startTime; + if (cost <= 1000 && !isError) { + referenceSum.setOneSecondLess(1L); + } else if (1000 < cost && cost <= 3000 && !isError) { + referenceSum.setThreeSecondLess(1L); + } else if (3000 < cost && cost <= 5000 && !isError) { + referenceSum.setFiveSecondLess(1L); + } else if (5000 < cost && !isError) { + referenceSum.setFiveSecondGreater(1L); + } else { + referenceSum.setError(1L); + } + referenceSum.setSummary(1L); + return referenceSum; + } + + @Override public void parseFirst(SpanObject spanObject, int applicationId, int applicationInstanceId) { + timeBucket = TimeBucketUtils.INSTANCE.getMinuteTimeBucket(spanObject.getStartTime()); + } + + @Override public void parseRef(TraceSegmentReference reference, int applicationId, int applicationInstanceId) { + hasReference = true; + } + + @Override public void build() { + logger.debug("node reference summary listener build"); + StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); + if (!hasReference) { + nodeExitReferences.addAll(nodeEntryReferences); + } + + for (NodeRefSumDataDefine.NodeReferenceSum referenceSum : nodeExitReferences) { + referenceSum.setId(timeBucket + Const.ID_SPLIT + referenceSum.getAgg()); + referenceSum.setTimeBucket(timeBucket); + + try { + logger.debug("send to node reference summary aggregation worker, id: {}", referenceSum.getId()); + context.getClusterWorkerContext().lookup(NodeRefSumAggregationWorker.WorkerRole.INSTANCE).tell(referenceSum.transform()); + } 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/noderef/summary/dao/INodeRefSumDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/INodeRefSumDAO.java new file mode 100644 index 000000000..ef4f69737 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/INodeRefSumDAO.java @@ -0,0 +1,12 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.summary.dao; + +import java.util.List; +import java.util.Map; +import org.skywalking.apm.collector.stream.worker.impl.data.Data; + +/** + * @author pengys5 + */ +public interface INodeRefSumDAO { + List prepareBatch(Map dataMap); +} 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 new file mode 100644 index 000000000..a20099784 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumEsDAO.java @@ -0,0 +1,35 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.summary.dao; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import org.elasticsearch.action.index.IndexRequestBuilder; +import org.skywalking.apm.collector.agentstream.worker.noderef.summary.define.NodeRefSumTable; +import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.stream.worker.impl.data.Data; + +/** + * @author pengys5 + */ +public class NodeRefSumEsDAO extends EsDAO implements INodeRefSumDAO { + + @Override public List prepareBatch(Map dataMap) { + List indexRequestBuilders = new ArrayList<>(); + dataMap.forEach((id, data) -> { + Map source = new HashMap(); + source.put(NodeRefSumTable.COLUMN_ONE_SECOND_LESS, data.getDataLong(0)); + source.put(NodeRefSumTable.COLUMN_THREE_SECOND_LESS, data.getDataLong(1)); + source.put(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS, data.getDataLong(2)); + source.put(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER, data.getDataLong(3)); + source.put(NodeRefSumTable.COLUMN_ERROR, data.getDataLong(4)); + source.put(NodeRefSumTable.COLUMN_SUMMARY, data.getDataLong(5)); + source.put(NodeRefSumTable.COLUMN_AGG, data.getDataString(1)); + source.put(NodeRefSumTable.COLUMN_TIME_BUCKET, data.getDataLong(6)); + + IndexRequestBuilder builder = getClient().prepareIndex(NodeRefSumTable.TABLE, id).setSource(source); + indexRequestBuilders.add(builder); + }); + return indexRequestBuilders; + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumH2DAO.java new file mode 100644 index 000000000..2a3b83501 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/dao/NodeRefSumH2DAO.java @@ -0,0 +1,15 @@ +package org.skywalking.apm.collector.agentstream.worker.noderef.summary.dao; + +import java.util.List; +import java.util.Map; +import org.skywalking.apm.collector.storage.h2.dao.H2DAO; + +/** + * @author pengys5 + */ +public class NodeRefSumH2DAO extends H2DAO implements INodeRefSumDAO { + + @Override public List prepareBatch(Map map) { + return null; + } +} 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 758ecc16a..cecfae5b5 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 @@ -3,8 +3,10 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.summary.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.operate.CoverOperation; +import org.skywalking.apm.collector.stream.worker.impl.data.TransformToData; +import org.skywalking.apm.collector.stream.worker.impl.data.operate.AddOperation; import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation; /** @@ -24,14 +26,14 @@ public class NodeRefSumDataDefine extends DataDefine { @Override protected void attributeDefine() { addAttribute(0, new Attribute(NodeRefSumTable.COLUMN_ID, AttributeType.STRING, new NonOperation())); - addAttribute(1, new Attribute(NodeRefSumTable.COLUMN_ONE_SECOND_LESS, AttributeType.LONG, new NonOperation())); - addAttribute(2, new Attribute(NodeRefSumTable.COLUMN_THREE_SECOND_LESS, AttributeType.LONG, new NonOperation())); - addAttribute(3, new Attribute(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS, AttributeType.LONG, new NonOperation())); - addAttribute(4, new Attribute(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER, AttributeType.LONG, new NonOperation())); - addAttribute(5, new Attribute(NodeRefSumTable.COLUMN_ERROR, AttributeType.LONG, new NonOperation())); - addAttribute(6, new Attribute(NodeRefSumTable.COLUMN_SUMMARY, AttributeType.LONG, new NonOperation())); - addAttribute(7, new Attribute(NodeRefSumTable.COLUMN_AGG, AttributeType.STRING, new CoverOperation())); - addAttribute(8, new Attribute(NodeRefSumTable.COLUMN_TIME_BUCKET, AttributeType.LONG, new CoverOperation())); + addAttribute(1, new Attribute(NodeRefSumTable.COLUMN_ONE_SECOND_LESS, AttributeType.LONG, new AddOperation())); + addAttribute(2, new Attribute(NodeRefSumTable.COLUMN_THREE_SECOND_LESS, AttributeType.LONG, new AddOperation())); + addAttribute(3, new Attribute(NodeRefSumTable.COLUMN_FIVE_SECOND_LESS, AttributeType.LONG, new AddOperation())); + addAttribute(4, new Attribute(NodeRefSumTable.COLUMN_FIVE_SECOND_GREATER, AttributeType.LONG, new AddOperation())); + addAttribute(5, new Attribute(NodeRefSumTable.COLUMN_ERROR, AttributeType.LONG, new AddOperation())); + addAttribute(6, new Attribute(NodeRefSumTable.COLUMN_SUMMARY, AttributeType.LONG, new AddOperation())); + addAttribute(7, new Attribute(NodeRefSumTable.COLUMN_AGG, AttributeType.STRING, new NonOperation())); + addAttribute(8, new Attribute(NodeRefSumTable.COLUMN_TIME_BUCKET, AttributeType.LONG, new NonOperation())); } @Override public Object deserialize(RemoteData remoteData) { @@ -62,14 +64,14 @@ public class NodeRefSumDataDefine extends DataDefine { return builder.build(); } - public static class NodeReferenceSum { + public static class NodeReferenceSum implements TransformToData { private String id; - private Long oneSecondLess; - private Long threeSecondLess; - private Long fiveSecondLess; - private Long fiveSecondGreater; - private Long error; - private Long summary; + private Long oneSecondLess = 0L; + private Long threeSecondLess = 0L; + private Long fiveSecondLess = 0L; + private Long fiveSecondGreater = 0L; + private Long error = 0L; + private Long summary = 0L; private String agg; private long timeBucket; @@ -89,6 +91,21 @@ public class NodeRefSumDataDefine extends DataDefine { public NodeReferenceSum() { } + @Override public Data transform() { + NodeRefSumDataDefine define = new NodeRefSumDataDefine(); + Data data = define.build(id); + data.setDataString(0, this.id); + data.setDataString(1, this.agg); + data.setDataLong(0, this.oneSecondLess); + data.setDataLong(1, this.threeSecondLess); + data.setDataLong(2, this.fiveSecondLess); + data.setDataLong(3, this.fiveSecondGreater); + data.setDataLong(4, this.error); + data.setDataLong(5, this.summary); + data.setDataLong(6, this.timeBucket); + return data; + } + public String getId() { return id; } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumEsTableDefine.java index c10b0d4a6..6fd0b818b 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/summary/define/NodeRefSumEsTableDefine.java @@ -13,7 +13,7 @@ public class NodeRefSumEsTableDefine extends ElasticSearchTableDefine { } @Override public int refreshInterval() { - return 0; + return 2; } @Override public int numberOfShards() { 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 701109755..64bc47399 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 @@ -5,6 +5,7 @@ import java.util.List; import org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentSpanListener; import org.skywalking.apm.collector.agentstream.worker.node.mapping.NodeMappingSpanListener; import org.skywalking.apm.collector.agentstream.worker.noderef.reference.NodeRefSpanListener; +import org.skywalking.apm.collector.agentstream.worker.noderef.summary.NodeRefSumSpanListener; import org.skywalking.apm.collector.agentstream.worker.segment.define.SegmentDataDefine; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.core.util.CollectionUtils; @@ -36,10 +37,12 @@ public class SegmentParse { spanListeners.add(new NodeComponentSpanListener()); spanListeners.add(new NodeMappingSpanListener()); spanListeners.add(new NodeRefSpanListener()); + spanListeners.add(new NodeRefSumSpanListener()); refsListeners = new ArrayList<>(); refsListeners.add(new NodeMappingSpanListener()); refsListeners.add(new NodeRefSpanListener()); + refsListeners.add(new NodeRefSumSpanListener()); } public void parse(List traceIds, TraceSegmentObject segmentObject) { 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 0eff3bcc1..5bd0ff7b1 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 @@ -4,4 +4,5 @@ org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.Service org.skywalking.apm.collector.agentstream.worker.node.component.dao.NodeComponentEsDAO org.skywalking.apm.collector.agentstream.worker.node.mapping.dao.NodeMappingEsDAO org.skywalking.apm.collector.agentstream.worker.noderef.reference.dao.NodeReferenceEsDAO -org.skywalking.apm.collector.agentstream.worker.segment.dao.SegmentEsDAO \ No newline at end of file +org.skywalking.apm.collector.agentstream.worker.segment.dao.SegmentEsDAO +org.skywalking.apm.collector.agentstream.worker.noderef.summary.dao.NodeRefSumEsDAO \ 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 0d5afd43f..8e740963d 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 @@ -4,4 +4,5 @@ org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.Service org.skywalking.apm.collector.agentstream.worker.node.component.dao.NodeComponentH2DAO org.skywalking.apm.collector.agentstream.worker.node.mapping.dao.NodeMappingH2DAO org.skywalking.apm.collector.agentstream.worker.noderef.reference.dao.NodeReferenceH2DAO -org.skywalking.apm.collector.agentstream.worker.segment.dao.SegmentH2DAO \ No newline at end of file +org.skywalking.apm.collector.agentstream.worker.segment.dao.SegmentH2DAO +org.skywalking.apm.collector.agentstream.worker.noderef.summary.dao.NodeRefSumH2DAO \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_async_worker_provider.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_async_worker_provider.define index 389c3a19d..2e72aac3c 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_async_worker_provider.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_async_worker_provider.define @@ -7,6 +7,9 @@ org.skywalking.apm.collector.agentstream.worker.node.mapping.NodeMappingPersiste org.skywalking.apm.collector.agentstream.worker.noderef.reference.NodeRefAggregationWorker$Factory org.skywalking.apm.collector.agentstream.worker.noderef.reference.NodeRefPersistenceWorker$Factory +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.segment.SegmentPersistenceWorker$Factory org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationRegisterSerialWorker$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 ae963c8b8..b82aef6bd 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 @@ -4,4 +4,5 @@ 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 \ No newline at end of file +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 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 2bb31fb38..4b726cba3 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 @@ -7,6 +7,9 @@ org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeMappingH org.skywalking.apm.collector.agentstream.worker.noderef.reference.define.NodeRefEsTableDefine org.skywalking.apm.collector.agentstream.worker.noderef.reference.define.NodeRefH2TableDefine +org.skywalking.apm.collector.agentstream.worker.noderef.summary.define.NodeRefSumEsTableDefine +org.skywalking.apm.collector.agentstream.worker.noderef.summary.define.NodeRefSumH2TableDefine + org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationEsTableDefine org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationH2TableDefine diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/operate/AddOperation.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/operate/AddOperation.java new file mode 100644 index 000000000..d625b279a --- /dev/null +++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/impl/data/operate/AddOperation.java @@ -0,0 +1,29 @@ +package org.skywalking.apm.collector.stream.worker.impl.data.operate; + +import org.skywalking.apm.collector.stream.worker.impl.data.Operation; + +/** + * @author pengys5 + */ +public class AddOperation implements Operation { + + @Override public String operate(String newValue, String oldValue) { + throw new UnsupportedOperationException("not support string addition operation"); + } + + @Override public Long operate(Long newValue, Long oldValue) { + return newValue + oldValue; + } + + @Override public Float operate(Float newValue, Float oldValue) { + return newValue + oldValue; + } + + @Override public Integer operate(Integer newValue, Integer oldValue) { + return newValue + oldValue; + } + + @Override public byte[] operate(byte[] newValue, byte[] oldValue) { + throw new UnsupportedOperationException("not support byte addition operation"); + } +}