node reference summary save to es success
This commit is contained in:
parent
e216a369b3
commit
430fe745a4
|
|
@ -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<NodeRefSumAggregationWorker> {
|
||||
@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();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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<String, Data> dataMap) {
|
||||
INodeRefSumDAO dao = (INodeRefSumDAO)DAOContainer.INSTANCE.get(INodeRefSumDAO.class.getName());
|
||||
return dao.prepareBatch(dataMap);
|
||||
}
|
||||
|
||||
public static class Factory extends AbstractLocalAsyncWorkerProvider<NodeRefSumPersistenceWorker> {
|
||||
@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();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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<NodeRefSumRemoteWorker> {
|
||||
@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();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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<NodeRefSumDataDefine.NodeReferenceSum> nodeExitReferences = new ArrayList<>();
|
||||
private List<NodeRefSumDataDefine.NodeReferenceSum> 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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<String, Data> dataMap);
|
||||
}
|
||||
|
|
@ -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<String, Data> dataMap) {
|
||||
List<IndexRequestBuilder> indexRequestBuilders = new ArrayList<>();
|
||||
dataMap.forEach((id, data) -> {
|
||||
Map<String, Object> 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;
|
||||
}
|
||||
}
|
||||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
@ -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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -13,7 +13,7 @@ public class NodeRefSumEsTableDefine extends ElasticSearchTableDefine {
|
|||
}
|
||||
|
||||
@Override public int refreshInterval() {
|
||||
return 0;
|
||||
return 2;
|
||||
}
|
||||
|
||||
@Override public int numberOfShards() {
|
||||
|
|
|
|||
|
|
@ -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<UniqueId> traceIds, TraceSegmentObject segmentObject) {
|
||||
|
|
|
|||
|
|
@ -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
|
||||
org.skywalking.apm.collector.agentstream.worker.segment.dao.SegmentEsDAO
|
||||
org.skywalking.apm.collector.agentstream.worker.noderef.summary.dao.NodeRefSumEsDAO
|
||||
|
|
@ -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
|
||||
org.skywalking.apm.collector.agentstream.worker.segment.dao.SegmentH2DAO
|
||||
org.skywalking.apm.collector.agentstream.worker.noderef.summary.dao.NodeRefSumH2DAO
|
||||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
org.skywalking.apm.collector.agentstream.worker.noderef.reference.NodeRefRemoteWorker$Factory
|
||||
org.skywalking.apm.collector.agentstream.worker.noderef.summary.NodeRefSumRemoteWorker$Factory
|
||||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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");
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue