Merge branch 'master' into zhangxin/feature/add-default-span-layer
This commit is contained in:
commit
d9ba6e8221
|
|
@ -16,7 +16,7 @@ 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.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.stream.worker.util.TimeBucketUtils;
|
||||
import org.skywalking.apm.network.proto.CPU;
|
||||
import org.skywalking.apm.network.proto.Downstream;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentjvm.worker.cpu.define;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentjvm.worker.gc.define;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentjvm.worker.memory.define;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentjvm.worker.memorypool.define;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.cache;
|
|||
|
||||
import com.google.common.cache.Cache;
|
||||
import com.google.common.cache.CacheBuilder;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.IServiceNameDAO;
|
||||
import org.skywalking.apm.collector.storage.dao.DAOContainer;
|
||||
|
||||
|
|
|
|||
|
|
@ -4,8 +4,9 @@ import java.util.HashMap;
|
|||
import java.util.Map;
|
||||
import org.elasticsearch.action.index.IndexRequestBuilder;
|
||||
import org.elasticsearch.action.update.UpdateRequestBuilder;
|
||||
import org.skywalking.apm.collector.agentstream.worker.global.define.GlobalTraceTable;
|
||||
import org.skywalking.apm.collector.core.framework.UnexpectedException;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.skywalking.apm.collector.storage.table.global.GlobalTraceTable;
|
||||
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;
|
||||
|
|
@ -20,11 +21,11 @@ public class GlobalTraceEsDAO extends EsDAO implements IGlobalTraceDAO, IPersist
|
|||
private final Logger logger = LoggerFactory.getLogger(GlobalTraceEsDAO.class);
|
||||
|
||||
@Override public Data get(String id, DataDefine dataDefine) {
|
||||
return null;
|
||||
throw new UnexpectedException("There is no need to merge stream data with database data.");
|
||||
}
|
||||
|
||||
@Override public UpdateRequestBuilder prepareBatchUpdate(Data data) {
|
||||
return null;
|
||||
throw new UnexpectedException("There is no need to merge stream data with database data.");
|
||||
}
|
||||
|
||||
@Override public IndexRequestBuilder prepareBatchInsert(Data data) {
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.global.define;
|
||||
|
||||
import org.skywalking.apm.collector.remote.grpc.proto.RemoteData;
|
||||
import org.skywalking.apm.collector.storage.table.global.GlobalTraceTable;
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.global.define;
|
|||
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.global.GlobalTraceTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.global.define;
|
|||
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2TableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.global.GlobalTraceTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.node.component;
|
|||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener;
|
||||
|
|
|
|||
|
|
@ -5,7 +5,7 @@ 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.node.component.define.NodeComponentTable;
|
||||
import org.skywalking.apm.collector.storage.table.node.NodeComponentTable;
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.node.component.define;
|
||||
|
||||
import org.skywalking.apm.collector.remote.grpc.proto.RemoteData;
|
||||
import org.skywalking.apm.collector.storage.table.node.NodeComponentTable;
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.node.component.define;
|
|||
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.node.NodeComponentTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.node.component.define;
|
|||
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2TableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.node.NodeComponentTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,10 +0,0 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.node.component.define;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeComponentTable extends CommonTable {
|
||||
public static final String TABLE = "node_component";
|
||||
}
|
||||
|
|
@ -2,7 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.node.mapping;
|
|||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeMappingDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.FirstSpanListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.RefsListener;
|
||||
|
|
|
|||
|
|
@ -5,7 +5,7 @@ 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.node.mapping.define.NodeMappingTable;
|
||||
import org.skywalking.apm.collector.storage.table.node.NodeMappingTable;
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.node.mapping.define;
|
||||
|
||||
import org.skywalking.apm.collector.remote.grpc.proto.RemoteData;
|
||||
import org.skywalking.apm.collector.storage.table.node.NodeMappingTable;
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.node.mapping.define;
|
|||
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.node.NodeMappingTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.node.mapping.define;
|
|||
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2TableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.node.NodeMappingTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,10 +0,0 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.node.mapping.define;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeMappingTable extends CommonTable {
|
||||
public static final String TABLE = "node_mapping";
|
||||
}
|
||||
|
|
@ -2,7 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.reference;
|
|||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.cache.InstanceCache;
|
||||
import org.skywalking.apm.collector.agentstream.worker.noderef.reference.define.NodeRefDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener;
|
||||
|
|
|
|||
|
|
@ -5,7 +5,7 @@ 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.noderef.reference.define.NodeRefTable;
|
||||
import org.skywalking.apm.collector.storage.table.noderef.NodeRefTable;
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.noderef.reference.define;
|
||||
|
||||
import org.skywalking.apm.collector.remote.grpc.proto.RemoteData;
|
||||
import org.skywalking.apm.collector.storage.table.noderef.NodeRefTable;
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.reference.define
|
|||
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.noderef.NodeRefTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.reference.define
|
|||
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2TableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.noderef.NodeRefTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,10 +0,0 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.noderef.reference.define;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeRefTable extends CommonTable {
|
||||
public static final String TABLE = "node_reference";
|
||||
}
|
||||
|
|
@ -2,7 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.summary;
|
|||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.cache.InstanceCache;
|
||||
import org.skywalking.apm.collector.agentstream.worker.noderef.summary.define.NodeRefSumDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener;
|
||||
|
|
|
|||
|
|
@ -5,8 +5,8 @@ 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.noderef.reference.define.NodeRefTable;
|
||||
import org.skywalking.apm.collector.agentstream.worker.noderef.summary.define.NodeRefSumTable;
|
||||
import org.skywalking.apm.collector.storage.table.noderef.NodeRefTable;
|
||||
import org.skywalking.apm.collector.storage.table.noderef.NodeRefSumTable;
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.noderef.summary.define;
|
||||
|
||||
import org.skywalking.apm.collector.remote.grpc.proto.RemoteData;
|
||||
import org.skywalking.apm.collector.storage.table.noderef.NodeRefSumTable;
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.summary.define;
|
|||
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.noderef.NodeRefSumTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.summary.define;
|
|||
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2TableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.noderef.NodeRefSumTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.register.application;
|
||||
|
||||
import org.skywalking.apm.collector.remote.grpc.proto.RemoteData;
|
||||
import org.skywalking.apm.collector.storage.table.register.ApplicationTable;
|
||||
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.DataDefine;
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.register.application;
|
|||
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.register.ApplicationTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.register.application;
|
|||
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2TableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.register.ApplicationTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.register.application;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.IdAutoIncrement;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.application.dao.IApplicationDAO;
|
||||
import org.skywalking.apm.collector.storage.dao.DAOContainer;
|
||||
|
|
|
|||
|
|
@ -10,7 +10,7 @@ import org.elasticsearch.action.support.WriteRequest;
|
|||
import org.elasticsearch.index.query.QueryBuilders;
|
||||
import org.elasticsearch.search.SearchHit;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationTable;
|
||||
import org.skywalking.apm.collector.storage.table.register.ApplicationTable;
|
||||
import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.slf4j.Logger;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.register.instance;
|
||||
|
||||
import org.skywalking.apm.collector.remote.grpc.proto.RemoteData;
|
||||
import org.skywalking.apm.collector.storage.table.register.InstanceTable;
|
||||
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.DataDefine;
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.register.instance;
|
|||
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.register.InstanceTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.register.instance;
|
|||
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2TableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.register.InstanceTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -13,7 +13,7 @@ import org.elasticsearch.index.query.BoolQueryBuilder;
|
|||
import org.elasticsearch.index.query.QueryBuilders;
|
||||
import org.elasticsearch.search.SearchHit;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceTable;
|
||||
import org.skywalking.apm.collector.storage.table.register.InstanceTable;
|
||||
import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.slf4j.Logger;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.register.servicename;
|
||||
|
||||
import org.skywalking.apm.collector.remote.grpc.proto.RemoteData;
|
||||
import org.skywalking.apm.collector.storage.table.register.ServiceNameTable;
|
||||
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.DataDefine;
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.register.servicename;
|
|||
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.register.ServiceNameTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.register.servicename;
|
|||
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2TableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.register.ServiceNameTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -11,7 +11,7 @@ import org.elasticsearch.index.query.BoolQueryBuilder;
|
|||
import org.elasticsearch.index.query.QueryBuilders;
|
||||
import org.elasticsearch.search.SearchHit;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.servicename.ServiceNameDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.servicename.ServiceNameTable;
|
||||
import org.skywalking.apm.collector.storage.table.register.ServiceNameTable;
|
||||
import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.slf4j.Logger;
|
||||
|
|
|
|||
|
|
@ -4,7 +4,7 @@ import java.util.HashMap;
|
|||
import java.util.Map;
|
||||
import org.elasticsearch.action.index.IndexRequestBuilder;
|
||||
import org.elasticsearch.action.update.UpdateRequestBuilder;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.cost.define.SegmentCostTable;
|
||||
import org.skywalking.apm.collector.storage.table.segment.SegmentCostTable;
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.segment.cost.define;
|
||||
|
||||
import org.skywalking.apm.collector.remote.grpc.proto.RemoteData;
|
||||
import org.skywalking.apm.collector.storage.table.segment.SegmentCostTable;
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.segment.cost.define;
|
|||
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.segment.SegmentCostTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.segment.cost.define;
|
|||
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2TableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.segment.SegmentCostTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -5,7 +5,7 @@ import java.util.HashMap;
|
|||
import java.util.Map;
|
||||
import org.elasticsearch.action.index.IndexRequestBuilder;
|
||||
import org.elasticsearch.action.update.UpdateRequestBuilder;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.origin.define.SegmentTable;
|
||||
import org.skywalking.apm.collector.storage.table.segment.SegmentTable;
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.segment.origin.define;
|
|||
|
||||
import com.google.protobuf.ByteString;
|
||||
import org.skywalking.apm.collector.remote.grpc.proto.RemoteData;
|
||||
import org.skywalking.apm.collector.storage.table.segment.SegmentTable;
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.segment.origin.define;
|
|||
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.segment.SegmentTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.segment.origin.define;
|
|||
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine;
|
||||
import org.skywalking.apm.collector.storage.h2.define.H2TableDefine;
|
||||
import org.skywalking.apm.collector.storage.table.segment.SegmentTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.service.entry;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.core.util.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;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.service.entry.define;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.serviceref.reference;
|
|||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.cache.InstanceCache;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.EntrySpanListener;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.ExitSpanListener;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.serviceref.reference.define;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -3,6 +3,8 @@ package org.skywalking.apm.collector.agentstream.mock;
|
|||
import com.google.gson.JsonElement;
|
||||
import java.io.IOException;
|
||||
import org.skywalking.apm.collector.agentstream.HttpClientTools;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.application.dao.ApplicationEsDAO;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceDataDefine;
|
||||
import org.skywalking.apm.collector.agentstream.worker.register.instance.dao.InstanceEsDAO;
|
||||
import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
|
||||
|
|
@ -25,6 +27,14 @@ public class SegmentPost {
|
|||
InstanceDataDefine.Instance providerInstance = new InstanceDataDefine.Instance("3", 3, "dubbox-provider", 1501858094526L, 3);
|
||||
instanceEsDAO.save(providerInstance);
|
||||
|
||||
ApplicationEsDAO applicationEsDAO = new ApplicationEsDAO();
|
||||
applicationEsDAO.setClient(client);
|
||||
|
||||
ApplicationDataDefine.Application consumerApplication = new ApplicationDataDefine.Application("2", "dubbox-consumer", 2);
|
||||
applicationEsDAO.save(consumerApplication);
|
||||
ApplicationDataDefine.Application providerApplication = new ApplicationDataDefine.Application("3", "dubbox-provider", 3);
|
||||
applicationEsDAO.save(providerApplication);
|
||||
|
||||
JsonElement consumer = JsonFileReader.INSTANCE.read("json/segment/normal/dubbox-consumer.json");
|
||||
HttpClientTools.INSTANCE.post("http://localhost:12800/segments", consumer.toString());
|
||||
|
||||
|
|
|
|||
|
|
@ -23,11 +23,11 @@ public abstract class StorageInstaller {
|
|||
if (!isExists(client, tableDefine)) {
|
||||
logger.info("table: {} not exists", tableDefine.getName());
|
||||
tableDefine.initialize();
|
||||
createTable(client, tableDefine);
|
||||
} else {
|
||||
logger.info("table: {} exists", tableDefine.getName());
|
||||
deleteTable(client, tableDefine);
|
||||
// deleteTable(client, tableDefine);
|
||||
}
|
||||
createTable(client, tableDefine);
|
||||
}
|
||||
} catch (DefineException e) {
|
||||
throw new StorageInstallException(e.getMessage(), e);
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.stream.worker.util;
|
||||
package org.skywalking.apm.collector.core.util;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -9,4 +9,6 @@ public class Const {
|
|||
public static final int USER_ID = 1;
|
||||
public static final String USER_CODE = "User";
|
||||
public static final String SEGMENT_SPAN_SPLIT = "S";
|
||||
public static final String UNKNOWN = "Unknown";
|
||||
public static final String EXCEPTION = "Exception";
|
||||
}
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package org.skywalking.apm.collector.stream.worker.storage;
|
||||
package org.skywalking.apm.collector.storage.table;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.global.define;
|
||||
package org.skywalking.apm.collector.storage.table.global;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -0,0 +1,10 @@
|
|||
package org.skywalking.apm.collector.storage.table.node;
|
||||
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeComponentTable extends CommonTable {
|
||||
public static final String TABLE = "node_component";
|
||||
}
|
||||
|
|
@ -0,0 +1,10 @@
|
|||
package org.skywalking.apm.collector.storage.table.node;
|
||||
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeMappingTable extends CommonTable {
|
||||
public static final String TABLE = "node_mapping";
|
||||
}
|
||||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.noderef.summary.define;
|
||||
package org.skywalking.apm.collector.storage.table.noderef;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -0,0 +1,10 @@
|
|||
package org.skywalking.apm.collector.storage.table.noderef;
|
||||
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeRefTable extends CommonTable {
|
||||
public static final String TABLE = "node_reference";
|
||||
}
|
||||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.register.application;
|
||||
package org.skywalking.apm.collector.storage.table.register;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.register.instance;
|
||||
package org.skywalking.apm.collector.storage.table.register;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.register.servicename;
|
||||
package org.skywalking.apm.collector.storage.table.register;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.segment.cost.define;
|
||||
package org.skywalking.apm.collector.storage.table.segment;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -1,6 +1,6 @@
|
|||
package org.skywalking.apm.collector.agentstream.worker.segment.origin.define;
|
||||
package org.skywalking.apm.collector.storage.table.segment;
|
||||
|
||||
import org.skywalking.apm.collector.stream.worker.storage.CommonTable;
|
||||
import org.skywalking.apm.collector.storage.table.CommonTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -33,10 +33,5 @@
|
|||
<artifactId>apm-collector-storage</artifactId>
|
||||
<version>${project.version}</version>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.skywalking</groupId>
|
||||
<artifactId>apm-collector-agentstream</artifactId>
|
||||
<version>${project.version}</version>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
</project>
|
||||
|
|
@ -0,0 +1,27 @@
|
|||
package org.skywalking.apm.collector.ui.cache;
|
||||
|
||||
import com.google.common.cache.Cache;
|
||||
import com.google.common.cache.CacheBuilder;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.storage.dao.DAOContainer;
|
||||
import org.skywalking.apm.collector.ui.dao.IApplicationDAO;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class ApplicationCache {
|
||||
|
||||
//TODO size configuration
|
||||
private static Cache<Integer, String> CACHE = CacheBuilder.newBuilder().maximumSize(1000).build();
|
||||
|
||||
public static String get(int applicationId) {
|
||||
try {
|
||||
return CACHE.get(applicationId, () -> {
|
||||
IApplicationDAO dao = (IApplicationDAO)DAOContainer.INSTANCE.get(IApplicationDAO.class.getName());
|
||||
return dao.getApplicationCode(applicationId);
|
||||
});
|
||||
} catch (Throwable e) {
|
||||
return Const.EXCEPTION;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,27 @@
|
|||
package org.skywalking.apm.collector.ui.cache;
|
||||
|
||||
import com.google.common.cache.Cache;
|
||||
import com.google.common.cache.CacheBuilder;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.storage.dao.DAOContainer;
|
||||
import org.skywalking.apm.collector.ui.dao.IServiceNameDAO;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class ServiceNameCache {
|
||||
|
||||
//TODO size configuration
|
||||
private static Cache<Integer, String> CACHE = CacheBuilder.newBuilder().maximumSize(1000).build();
|
||||
|
||||
public static String get(int serviceId) {
|
||||
try {
|
||||
return CACHE.get(serviceId, () -> {
|
||||
IServiceNameDAO dao = (IServiceNameDAO)DAOContainer.INSTANCE.get(IServiceNameDAO.class.getName());
|
||||
return dao.getServiceName(serviceId);
|
||||
});
|
||||
} catch (Throwable e) {
|
||||
return Const.EXCEPTION;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,30 @@
|
|||
package org.skywalking.apm.collector.ui.dao;
|
||||
|
||||
import org.elasticsearch.action.get.GetRequestBuilder;
|
||||
import org.elasticsearch.action.get.GetResponse;
|
||||
import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.skywalking.apm.collector.storage.table.register.ApplicationTable;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class ApplicationEsDAO extends EsDAO implements IApplicationDAO {
|
||||
|
||||
private final Logger logger = LoggerFactory.getLogger(ApplicationEsDAO.class);
|
||||
|
||||
@Override public String getApplicationCode(int applicationId) {
|
||||
logger.debug("get application code, applicationId: {}", applicationId);
|
||||
ElasticSearchClient client = getClient();
|
||||
GetRequestBuilder getRequestBuilder = client.prepareGet(ApplicationTable.TABLE, String.valueOf(applicationId));
|
||||
|
||||
GetResponse getResponse = getRequestBuilder.get();
|
||||
if (getResponse.isExists()) {
|
||||
return (String)getResponse.getSource().get(ApplicationTable.COLUMN_APPLICATION_CODE);
|
||||
}
|
||||
return Const.UNKNOWN;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,13 @@
|
|||
package org.skywalking.apm.collector.ui.dao;
|
||||
|
||||
import org.skywalking.apm.collector.storage.h2.dao.H2DAO;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class ApplicationH2DAO extends H2DAO implements IApplicationDAO {
|
||||
|
||||
@Override public String getApplicationCode(int applicationId) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
|
@ -7,8 +7,8 @@ import org.elasticsearch.action.search.SearchResponse;
|
|||
import org.elasticsearch.action.search.SearchType;
|
||||
import org.elasticsearch.index.query.QueryBuilders;
|
||||
import org.elasticsearch.search.SearchHit;
|
||||
import org.skywalking.apm.collector.agentstream.worker.global.define.GlobalTraceTable;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.skywalking.apm.collector.storage.table.global.GlobalTraceTable;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
|
|
|
|||
|
|
@ -0,0 +1,8 @@
|
|||
package org.skywalking.apm.collector.ui.dao;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public interface IApplicationDAO {
|
||||
String getApplicationCode(int applicationId);
|
||||
}
|
||||
|
|
@ -0,0 +1,8 @@
|
|||
package org.skywalking.apm.collector.ui.dao;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public interface IServiceNameDAO {
|
||||
String getServiceName(int serviceId);
|
||||
}
|
||||
|
|
@ -8,9 +8,9 @@ import org.elasticsearch.action.search.SearchType;
|
|||
import org.elasticsearch.index.query.QueryBuilders;
|
||||
import org.elasticsearch.search.aggregations.AggregationBuilders;
|
||||
import org.elasticsearch.search.aggregations.bucket.terms.Terms;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.node.component.define.NodeComponentTable;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.skywalking.apm.collector.storage.table.node.NodeComponentTable;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
|
|
|
|||
|
|
@ -8,9 +8,9 @@ import org.elasticsearch.action.search.SearchType;
|
|||
import org.elasticsearch.index.query.QueryBuilders;
|
||||
import org.elasticsearch.search.aggregations.AggregationBuilders;
|
||||
import org.elasticsearch.search.aggregations.bucket.terms.Terms;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeMappingTable;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.skywalking.apm.collector.storage.table.node.NodeMappingTable;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
|
|
|
|||
|
|
@ -10,9 +10,9 @@ import org.elasticsearch.search.aggregations.AggregationBuilders;
|
|||
import org.elasticsearch.search.aggregations.bucket.terms.Terms;
|
||||
import org.elasticsearch.search.aggregations.bucket.terms.TermsAggregationBuilder;
|
||||
import org.elasticsearch.search.aggregations.metrics.sum.Sum;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.noderef.summary.define.NodeRefSumTable;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.skywalking.apm.collector.storage.table.noderef.NodeRefSumTable;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
|
|
|
|||
|
|
@ -8,9 +8,9 @@ import org.elasticsearch.action.search.SearchType;
|
|||
import org.elasticsearch.index.query.QueryBuilders;
|
||||
import org.elasticsearch.search.aggregations.AggregationBuilders;
|
||||
import org.elasticsearch.search.aggregations.bucket.terms.Terms;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.agentstream.worker.noderef.reference.define.NodeRefTable;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.skywalking.apm.collector.storage.table.noderef.NodeRefTable;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
|
|
|
|||
|
|
@ -12,12 +12,12 @@ import org.elasticsearch.index.query.QueryBuilders;
|
|||
import org.elasticsearch.index.query.RangeQueryBuilder;
|
||||
import org.elasticsearch.search.SearchHit;
|
||||
import org.elasticsearch.search.sort.SortOrder;
|
||||
import org.skywalking.apm.collector.agentstream.worker.global.define.GlobalTraceTable;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.cost.define.SegmentCostTable;
|
||||
import org.skywalking.apm.collector.core.util.CollectionUtils;
|
||||
import org.skywalking.apm.collector.core.util.StringUtils;
|
||||
import org.skywalking.apm.collector.storage.dao.DAOContainer;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.skywalking.apm.collector.storage.table.global.GlobalTraceTable;
|
||||
import org.skywalking.apm.collector.storage.table.segment.SegmentCostTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
|
|||
|
|
@ -4,9 +4,9 @@ import com.google.protobuf.InvalidProtocolBufferException;
|
|||
import java.util.Base64;
|
||||
import java.util.Map;
|
||||
import org.elasticsearch.action.get.GetResponse;
|
||||
import org.skywalking.apm.collector.agentstream.worker.segment.origin.define.SegmentTable;
|
||||
import org.skywalking.apm.collector.core.util.StringUtils;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.skywalking.apm.collector.storage.table.segment.SegmentTable;
|
||||
import org.skywalking.apm.network.proto.TraceSegmentObject;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
|
|
|||
|
|
@ -0,0 +1,25 @@
|
|||
package org.skywalking.apm.collector.ui.dao;
|
||||
|
||||
import org.elasticsearch.action.get.GetRequestBuilder;
|
||||
import org.elasticsearch.action.get.GetResponse;
|
||||
import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
|
||||
import org.skywalking.apm.collector.storage.table.register.ServiceNameTable;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class ServiceNameEsDAO extends EsDAO implements IServiceNameDAO {
|
||||
|
||||
@Override public String getServiceName(int serviceId) {
|
||||
ElasticSearchClient client = getClient();
|
||||
GetRequestBuilder getRequestBuilder = client.prepareGet(ServiceNameTable.TABLE, String.valueOf(serviceId));
|
||||
|
||||
GetResponse getResponse = getRequestBuilder.get();
|
||||
if (getResponse.isExists()) {
|
||||
return (String)getResponse.getSource().get(ServiceNameTable.COLUMN_SERVICE_NAME);
|
||||
}
|
||||
return Const.UNKNOWN;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,13 @@
|
|||
package org.skywalking.apm.collector.ui.dao;
|
||||
|
||||
import org.skywalking.apm.collector.storage.h2.dao.H2DAO;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class ServiceNameH2DAO extends H2DAO implements IServiceNameDAO {
|
||||
|
||||
@Override public String getServiceName(int serviceId) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
|
@ -4,6 +4,7 @@ import com.google.gson.JsonArray;
|
|||
import com.google.gson.JsonObject;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.collector.storage.dao.DAOContainer;
|
||||
import org.skywalking.apm.collector.ui.cache.ServiceNameCache;
|
||||
import org.skywalking.apm.collector.ui.dao.ISegmentDAO;
|
||||
import org.skywalking.apm.network.proto.KeyWithStringValue;
|
||||
import org.skywalking.apm.network.proto.LogMessage;
|
||||
|
|
@ -23,7 +24,11 @@ public class SpanService {
|
|||
List<SpanObject> spans = segmentObject.getSpansList();
|
||||
for (SpanObject spanObject : spans) {
|
||||
if (spanId == spanObject.getSpanId()) {
|
||||
spanJson.addProperty("operationName", spanObject.getOperationName());
|
||||
String operationName = spanObject.getOperationName();
|
||||
if (spanObject.getOperationNameId() != 0) {
|
||||
operationName = ServiceNameCache.get(spanObject.getOperationNameId());
|
||||
}
|
||||
spanJson.addProperty("operationName", operationName);
|
||||
spanJson.addProperty("startTime", spanObject.getStartTime());
|
||||
spanJson.addProperty("endTime", spanObject.getEndTime());
|
||||
|
||||
|
|
|
|||
|
|
@ -4,7 +4,7 @@ import com.google.gson.JsonArray;
|
|||
import com.google.gson.JsonObject;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
|
|
|
|||
|
|
@ -4,10 +4,12 @@ import com.google.gson.JsonArray;
|
|||
import com.google.gson.JsonObject;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.collector.stream.worker.util.Const;
|
||||
import org.skywalking.apm.collector.core.util.CollectionUtils;
|
||||
import org.skywalking.apm.collector.core.util.Const;
|
||||
import org.skywalking.apm.collector.core.util.ObjectUtils;
|
||||
import org.skywalking.apm.collector.storage.dao.DAOContainer;
|
||||
import org.skywalking.apm.collector.ui.cache.ApplicationCache;
|
||||
import org.skywalking.apm.collector.ui.cache.ServiceNameCache;
|
||||
import org.skywalking.apm.collector.ui.dao.IGlobalTraceDAO;
|
||||
import org.skywalking.apm.collector.ui.dao.ISegmentDAO;
|
||||
import org.skywalking.apm.network.proto.SpanObject;
|
||||
|
|
@ -93,8 +95,12 @@ public class TraceStackService {
|
|||
String segmentSpanId = segmentId + Const.SEGMENT_SPAN_SPLIT + String.valueOf(spanId);
|
||||
String segmentParentSpanId = segmentId + Const.SEGMENT_SPAN_SPLIT + String.valueOf(parentSpanId);
|
||||
long startTime = spanObject.getStartTime();
|
||||
|
||||
String operationName = spanObject.getOperationName();
|
||||
String applicationCode = "Code" + String.valueOf(segment.getApplicationId());
|
||||
if (spanObject.getOperationNameId() != 0) {
|
||||
operationName = ServiceNameCache.get(spanObject.getOperationNameId());
|
||||
}
|
||||
String applicationCode = ApplicationCache.get(segment.getApplicationId());
|
||||
|
||||
long cost = spanObject.getEndTime() - spanObject.getStartTime();
|
||||
if (cost == 0) {
|
||||
|
|
|
|||
|
|
@ -4,4 +4,6 @@ org.skywalking.apm.collector.ui.dao.NodeReferenceEsDAO
|
|||
org.skywalking.apm.collector.ui.dao.NodeRefSumEsDAO
|
||||
org.skywalking.apm.collector.ui.dao.SegmentCostEsDAO
|
||||
org.skywalking.apm.collector.ui.dao.GlobalTraceEsDAO
|
||||
org.skywalking.apm.collector.ui.dao.SegmentEsDAO
|
||||
org.skywalking.apm.collector.ui.dao.SegmentEsDAO
|
||||
org.skywalking.apm.collector.ui.dao.ApplicationEsDAO
|
||||
org.skywalking.apm.collector.ui.dao.ServiceNameEsDAO
|
||||
|
|
@ -4,4 +4,6 @@ org.skywalking.apm.collector.ui.dao.NodeReferenceH2DAO
|
|||
org.skywalking.apm.collector.ui.dao.NodeRefSumH2DAO
|
||||
org.skywalking.apm.collector.ui.dao.SegmentCostH2DAO
|
||||
org.skywalking.apm.collector.ui.dao.GlobalTraceH2DAO
|
||||
org.skywalking.apm.collector.ui.dao.SegmentH2DAO
|
||||
org.skywalking.apm.collector.ui.dao.SegmentH2DAO
|
||||
org.skywalking.apm.collector.ui.dao.ApplicationH2DAO
|
||||
org.skywalking.apm.collector.ui.dao.ServiceNameH2DAO
|
||||
|
|
@ -3,7 +3,6 @@ package org.skywalking.apm.agent.core.context;
|
|||
import java.util.LinkedList;
|
||||
import java.util.List;
|
||||
import org.skywalking.apm.agent.core.boot.ServiceManager;
|
||||
import org.skywalking.apm.agent.core.conf.RemoteDownstreamConfig;
|
||||
import org.skywalking.apm.agent.core.context.trace.AbstractSpan;
|
||||
import org.skywalking.apm.agent.core.context.trace.AbstractTracingSpan;
|
||||
import org.skywalking.apm.agent.core.context.trace.EntrySpan;
|
||||
|
|
@ -293,7 +292,7 @@ public class TracingContext implements AbstractTracerContext {
|
|||
@Override
|
||||
public Object doProcess(final int peerId) {
|
||||
return DictionaryManager.findOperationNameCodeSection()
|
||||
.findOnly(RemoteDownstreamConfig.Agent.APPLICATION_ID, operationName)
|
||||
.findOnly(segment.getApplicationId(), operationName)
|
||||
.doInCondition(
|
||||
new PossibleFound.FoundAndObtain() {
|
||||
@Override
|
||||
|
|
@ -312,7 +311,7 @@ public class TracingContext implements AbstractTracerContext {
|
|||
@Override
|
||||
public Object doProcess() {
|
||||
return DictionaryManager.findOperationNameCodeSection()
|
||||
.findOnly(RemoteDownstreamConfig.Agent.APPLICATION_ID, operationName)
|
||||
.findOnly(segment.getApplicationId(), operationName)
|
||||
.doInCondition(
|
||||
new PossibleFound.FoundAndObtain() {
|
||||
@Override
|
||||
|
|
|
|||
Loading…
Reference in New Issue