Node mapping save to es success.
This commit is contained in:
parent
d80e2ff43d
commit
1ca8da9570
|
|
@ -0,0 +1,66 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.node.mapping;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeMappingDataDefine;
|
||||
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 NodeMappingAggregationWorker extends AggregationWorker {
|
||||
|
||||
public NodeMappingAggregationWorker(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(NodeMappingRemoteWorker.WorkerRole.INSTANCE);
|
||||
}
|
||||
|
||||
public static class Factory extends AbstractLocalAsyncWorkerProvider<NodeMappingAggregationWorker> {
|
||||
@Override
|
||||
public Role role() {
|
||||
return WorkerRole.INSTANCE;
|
||||
}
|
||||
|
||||
@Override
|
||||
public NodeMappingAggregationWorker workerInstance(ClusterWorkerContext clusterContext) {
|
||||
return new NodeMappingAggregationWorker(role(), clusterContext);
|
||||
}
|
||||
|
||||
@Override
|
||||
public int queueSize() {
|
||||
return 1024;
|
||||
}
|
||||
}
|
||||
|
||||
public enum WorkerRole implements Role {
|
||||
INSTANCE;
|
||||
|
||||
@Override
|
||||
public String roleName() {
|
||||
return NodeMappingAggregationWorker.class.getSimpleName();
|
||||
}
|
||||
|
||||
@Override
|
||||
public WorkerSelector workerSelector() {
|
||||
return new HashCodeSelector();
|
||||
}
|
||||
|
||||
@Override public DataDefine dataDefine() {
|
||||
return new NodeMappingDataDefine();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,70 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.node.mapping;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import org.skywalking.apm.collector.agentstream.worker.node.mapping.dao.INodeMappingDAO;
|
||||
import org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeMappingDataDefine;
|
||||
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 NodeMappingPersistenceWorker extends PersistenceWorker {
|
||||
|
||||
public NodeMappingPersistenceWorker(Role role, ClusterWorkerContext clusterContext) {
|
||||
super(role, clusterContext);
|
||||
}
|
||||
|
||||
@Override public void preStart() throws ProviderNotFoundException {
|
||||
super.preStart();
|
||||
}
|
||||
|
||||
@Override protected List<?> prepareBatch(Map<String, Data> dataMap) {
|
||||
INodeMappingDAO dao = (INodeMappingDAO)DAOContainer.INSTANCE.get(INodeMappingDAO.class.getName());
|
||||
return dao.prepareBatch(dataMap);
|
||||
}
|
||||
|
||||
public static class Factory extends AbstractLocalAsyncWorkerProvider<NodeMappingPersistenceWorker> {
|
||||
@Override
|
||||
public Role role() {
|
||||
return WorkerRole.INSTANCE;
|
||||
}
|
||||
|
||||
@Override
|
||||
public NodeMappingPersistenceWorker workerInstance(ClusterWorkerContext clusterContext) {
|
||||
return new NodeMappingPersistenceWorker(role(), clusterContext);
|
||||
}
|
||||
|
||||
@Override
|
||||
public int queueSize() {
|
||||
return 1024;
|
||||
}
|
||||
}
|
||||
|
||||
public enum WorkerRole implements Role {
|
||||
INSTANCE;
|
||||
|
||||
@Override
|
||||
public String roleName() {
|
||||
return NodeMappingPersistenceWorker.class.getSimpleName();
|
||||
}
|
||||
|
||||
@Override
|
||||
public WorkerSelector workerSelector() {
|
||||
return new HashCodeSelector();
|
||||
}
|
||||
|
||||
@Override public DataDefine dataDefine() {
|
||||
return new NodeMappingDataDefine();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,60 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.node.mapping;
|
||||
|
||||
import org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeMappingDataDefine;
|
||||
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 NodeMappingRemoteWorker extends AbstractRemoteWorker {
|
||||
|
||||
protected NodeMappingRemoteWorker(Role role, ClusterWorkerContext clusterContext) {
|
||||
super(role, clusterContext);
|
||||
}
|
||||
|
||||
@Override public void preStart() throws ProviderNotFoundException {
|
||||
|
||||
}
|
||||
|
||||
@Override protected void onWork(Object message) throws WorkerException {
|
||||
getClusterContext().lookup(NodeMappingPersistenceWorker.WorkerRole.INSTANCE).tell(message);
|
||||
}
|
||||
|
||||
public static class Factory extends AbstractRemoteWorkerProvider<NodeMappingRemoteWorker> {
|
||||
@Override
|
||||
public Role role() {
|
||||
return WorkerRole.INSTANCE;
|
||||
}
|
||||
|
||||
@Override
|
||||
public NodeMappingRemoteWorker workerInstance(ClusterWorkerContext clusterContext) {
|
||||
return new NodeMappingRemoteWorker(role(), clusterContext);
|
||||
}
|
||||
}
|
||||
|
||||
public enum WorkerRole implements Role {
|
||||
INSTANCE;
|
||||
|
||||
@Override
|
||||
public String roleName() {
|
||||
return NodeMappingRemoteWorker.class.getSimpleName();
|
||||
}
|
||||
|
||||
@Override
|
||||
public WorkerSelector workerSelector() {
|
||||
return new HashCodeSelector();
|
||||
}
|
||||
|
||||
@Override public DataDefine dataDefine() {
|
||||
return new NodeMappingDataDefine();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -7,6 +7,11 @@ import org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeM
|
|||
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;
|
||||
|
|
@ -23,6 +28,7 @@ public class NodeMappingSpanListener implements RefsListener, FirstSpanListener
|
|||
private long timeBucket;
|
||||
|
||||
@Override public void parseRef(TraceSegmentReference reference, int applicationId, int applicationInstanceId) {
|
||||
logger.debug("node mapping listener parse reference");
|
||||
String peers = Const.PEERS_FRONT_SPLIT + reference.getNetworkAddressId() + Const.PEERS_BEHIND_SPLIT;
|
||||
if (reference.getNetworkAddressId() == 0) {
|
||||
peers = Const.PEERS_FRONT_SPLIT + reference.getNetworkAddress() + Const.PEERS_BEHIND_SPLIT;
|
||||
|
|
@ -37,11 +43,20 @@ public class NodeMappingSpanListener implements RefsListener, FirstSpanListener
|
|||
}
|
||||
|
||||
@Override public void build() {
|
||||
logger.debug("node mapping listener build");
|
||||
StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME);
|
||||
for (String agg : nodeMappings) {
|
||||
NodeMappingDataDefine.NodeMapping nodeMapping = new NodeMappingDataDefine.NodeMapping();
|
||||
nodeMapping.setId(timeBucket + Const.ID_SPLIT + agg);
|
||||
nodeMapping.setAgg(agg);
|
||||
nodeMapping.setTimeBucket(timeBucket);
|
||||
|
||||
try {
|
||||
logger.debug("send to node mapping aggregation worker, id: {}", nodeMapping.getId());
|
||||
context.getClusterWorkerContext().lookup(NodeMappingAggregationWorker.WorkerRole.INSTANCE).tell(nodeMapping.transform());
|
||||
} catch (WorkerInvokeException | WorkerNotFoundException e) {
|
||||
logger.error(e.getMessage(), e);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,12 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.node.mapping.dao;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.data.Data;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public interface INodeMappingDAO {
|
||||
List<?> prepareBatch(Map<String, Data> dataMap);
|
||||
}
|
||||
|
|
@ -0,0 +1,29 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.node.mapping.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.node.mapping.define.NodeMappingTable;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.data.Data;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeMappingEsDAO extends EsDAO implements INodeMappingDAO {
|
||||
|
||||
@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(NodeMappingTable.COLUMN_AGG, data.getDataString(1));
|
||||
source.put(NodeMappingTable.COLUMN_TIME_BUCKET, data.getDataLong(0));
|
||||
|
||||
IndexRequestBuilder builder = getClient().prepareIndex(NodeMappingTable.TABLE, id).setSource();
|
||||
indexRequestBuilders.add(builder);
|
||||
});
|
||||
return indexRequestBuilders;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,15 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.node.mapping.dao;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import org.skywalking.apm.collector.storage.h2.dao.H2DAO;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeMappingH2DAO extends H2DAO implements INodeMappingDAO {
|
||||
|
||||
@Override public List<?> prepareBatch(Map map) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
|
@ -3,7 +3,9 @@ package org.skywalking.apm.collector.agentstream.worker.node.mapping.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.TransformToData;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.data.operate.CoverOperation;
|
||||
import org.skywalking.apm.collector.stream.worker.impl.data.operate.NonOperation;
|
||||
|
||||
|
|
@ -34,7 +36,7 @@ public class NodeMappingDataDefine extends DataDefine {
|
|||
return null;
|
||||
}
|
||||
|
||||
public static class NodeMapping {
|
||||
public static class NodeMapping implements TransformToData {
|
||||
private String id;
|
||||
private String agg;
|
||||
private long timeBucket;
|
||||
|
|
@ -48,6 +50,15 @@ public class NodeMappingDataDefine extends DataDefine {
|
|||
public NodeMapping() {
|
||||
}
|
||||
|
||||
@Override public Data transform() {
|
||||
NodeMappingDataDefine define = new NodeMappingDataDefine();
|
||||
Data data = define.build(id);
|
||||
data.setDataString(0, this.id);
|
||||
data.setDataString(1, this.agg);
|
||||
data.setDataLong(0, this.timeBucket);
|
||||
return data;
|
||||
}
|
||||
|
||||
public String getId() {
|
||||
return id;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -13,7 +13,7 @@ public class NodeMappingEsTableDefine extends ElasticSearchTableDefine {
|
|||
}
|
||||
|
||||
@Override public int refreshInterval() {
|
||||
return 0;
|
||||
return 2;
|
||||
}
|
||||
|
||||
@Override public int numberOfShards() {
|
||||
|
|
|
|||
|
|
@ -31,6 +31,7 @@ public class SegmentParse {
|
|||
spanListeners.add(new NodeMappingSpanListener());
|
||||
|
||||
refsListeners = new ArrayList<>();
|
||||
refsListeners.add(new NodeMappingSpanListener());
|
||||
}
|
||||
|
||||
public void parse(List<UniqueId> traceIds, TraceSegmentObject segmentObject) {
|
||||
|
|
|
|||
|
|
@ -1,4 +1,5 @@
|
|||
org.skywalking.apm.collector.agentstream.worker.register.application.dao.ApplicationEsDAO
|
||||
org.skywalking.apm.collector.agentstream.worker.register.instance.dao.InstanceEsDAO
|
||||
org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.ServiceNameEsDAO
|
||||
org.skywalking.apm.collector.agentstream.worker.node.component.dao.NodeComponentEsDAO
|
||||
org.skywalking.apm.collector.agentstream.worker.node.component.dao.NodeComponentEsDAO
|
||||
org.skywalking.apm.collector.agentstream.worker.node.mapping.dao.NodeMappingEsDAO
|
||||
|
|
@ -1,4 +1,5 @@
|
|||
org.skywalking.apm.collector.agentstream.worker.register.application.dao.ApplicationH2DAO
|
||||
org.skywalking.apm.collector.agentstream.worker.register.instance.dao.InstanceH2DAO
|
||||
org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.ServiceNameH2DAO
|
||||
org.skywalking.apm.collector.agentstream.worker.node.component.dao.NodeComponentH2DAO
|
||||
org.skywalking.apm.collector.agentstream.worker.node.component.dao.NodeComponentH2DAO
|
||||
org.skywalking.apm.collector.agentstream.worker.node.mapping.dao.NodeMappingH2DAO
|
||||
|
|
@ -1,5 +1,9 @@
|
|||
org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentAggregationWorker$Factory
|
||||
org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentPersistenceWorker$Factory
|
||||
|
||||
org.skywalking.apm.collector.agentstream.worker.node.mapping.NodeMappingAggregationWorker$Factory
|
||||
org.skywalking.apm.collector.agentstream.worker.node.mapping.NodeMappingPersistenceWorker$Factory
|
||||
|
||||
org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationRegisterSerialWorker$Factory
|
||||
org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceRegisterSerialWorker$Factory
|
||||
org.skywalking.apm.collector.agentstream.worker.register.servicename.ServiceNameRegisterSerialWorker$Factory
|
||||
|
|
@ -2,4 +2,5 @@ org.skywalking.apm.collector.agentstream.worker.register.application.Application
|
|||
org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceRegisterRemoteWorker$Factory
|
||||
org.skywalking.apm.collector.agentstream.worker.register.servicename.ServiceNameRegisterRemoteWorker$Factory
|
||||
|
||||
org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentRemoteWorker$Factory
|
||||
org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentRemoteWorker$Factory
|
||||
org.skywalking.apm.collector.agentstream.worker.node.mapping.NodeMappingRemoteWorker$Factory
|
||||
|
|
@ -9,6 +9,7 @@ import org.skywalking.apm.network.proto.SpanLayer;
|
|||
import org.skywalking.apm.network.proto.SpanObject;
|
||||
import org.skywalking.apm.network.proto.SpanType;
|
||||
import org.skywalking.apm.network.proto.TraceSegmentObject;
|
||||
import org.skywalking.apm.network.proto.TraceSegmentReference;
|
||||
import org.skywalking.apm.network.proto.TraceSegmentServiceGrpc;
|
||||
import org.skywalking.apm.network.proto.UniqueId;
|
||||
import org.skywalking.apm.network.proto.UpstreamSegment;
|
||||
|
|
@ -68,8 +69,8 @@ public class TraceSegmentServiceHandlerTestCase {
|
|||
long now = System.currentTimeMillis();
|
||||
|
||||
TraceSegmentObject.Builder segmentBuilder = TraceSegmentObject.newBuilder();
|
||||
segmentBuilder.setApplicationId(1);
|
||||
segmentBuilder.setApplicationInstanceId(1);
|
||||
segmentBuilder.setApplicationId(2);
|
||||
segmentBuilder.setApplicationInstanceId(2);
|
||||
segmentBuilder.setTraceSegmentId(UniqueId.newBuilder().addIdParts(200).addIdParts(200).addIdParts(200).build());
|
||||
|
||||
SpanObject.Builder span_0 = SpanObject.newBuilder();
|
||||
|
|
@ -83,10 +84,22 @@ public class TraceSegmentServiceHandlerTestCase {
|
|||
span_0.setComponentId(ComponentsDefine.TOMCAT.getId());
|
||||
span_0.setIsError(false);
|
||||
span_0.setSpanType(SpanType.Entry);
|
||||
span_0.setPeerId(0);
|
||||
span_0.setPeer("localhost:8080");
|
||||
span_0.setPeerId(2);
|
||||
span_0.setPeer("localhost:8082");
|
||||
segmentBuilder.addSpans(span_0);
|
||||
|
||||
TraceSegmentReference.Builder ref_0 = TraceSegmentReference.newBuilder();
|
||||
ref_0.setEntryServiceId(1);
|
||||
ref_0.setEntryServiceName("ServiceName");
|
||||
ref_0.setNetworkAddress("localhost:8081");
|
||||
ref_0.setNetworkAddressId(1);
|
||||
ref_0.setParentApplicationInstanceId(1);
|
||||
ref_0.setParentServiceId(1);
|
||||
ref_0.setParentServiceName("");
|
||||
ref_0.setParentSpanId(2);
|
||||
ref_0.setParentTraceSegmentId(UniqueId.newBuilder().addIdParts(100).addIdParts(100).addIdParts(100).build());
|
||||
segmentBuilder.addRefs(ref_0);
|
||||
|
||||
builder.setSegment(segmentBuilder.build().toByteString());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue