Service entry test success.
This commit is contained in:
parent
6856b896e1
commit
119d21d525
|
|
@ -22,13 +22,13 @@ public class NodeRefSumEsDAO extends EsDAO implements INodeRefSumDAO, IPersisten
|
|||
if (getResponse.isExists()) {
|
||||
Data data = dataDefine.build(id);
|
||||
Map<String, Object> 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 {
|
||||
|
|
|
|||
|
|
@ -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<UniqueId> traceIds, TraceSegmentObject segmentObject) {
|
||||
|
|
|
|||
|
|
@ -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<ServiceEntryAggregationWorker> {
|
||||
@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();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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<ServiceEntryPersistenceWorker> {
|
||||
@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();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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<ServiceEntryRemoteWorker> {
|
||||
@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();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.service.entry.dao;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public interface IServiceEntryDAO {
|
||||
}
|
||||
|
|
@ -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<IndexRequestBuilder, UpdateRequestBuilder> {
|
||||
|
||||
@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<String, Object> 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<String, Object> 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<String, Object> 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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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 {
|
||||
}
|
||||
|
|
@ -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<ServiceEntry> {
|
||||
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;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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()));
|
||||
}
|
||||
}
|
||||
|
|
@ -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()));
|
||||
}
|
||||
}
|
||||
|
|
@ -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";
|
||||
}
|
||||
|
|
@ -34,6 +34,7 @@ public class PersistenceTimer implements Starter {
|
|||
List<PersistenceWorker> 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();
|
||||
|
|
|
|||
|
|
@ -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
|
||||
org.skywalking.apm.collector.agentstream.worker.global.dao.GlobalTraceEsDAO
|
||||
org.skywalking.apm.collector.agentstream.worker.service.entry.dao.ServiceEntryEsDAO
|
||||
|
|
@ -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
|
||||
org.skywalking.apm.collector.agentstream.worker.global.dao.GlobalTraceH2DAO
|
||||
org.skywalking.apm.collector.agentstream.worker.service.entry.dao.ServiceEntryH2DAO
|
||||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
org.skywalking.apm.collector.agentstream.worker.noderef.summary.NodeRefSumRemoteWorker$Factory
|
||||
|
||||
org.skywalking.apm.collector.agentstream.worker.service.entry.ServiceEntryRemoteWorker$Factory
|
||||
|
|
@ -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
|
||||
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
|
||||
|
|
@ -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();
|
||||
|
||||
|
|
|
|||
Loading…
Reference in New Issue