diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandler.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandler.java index 15bc9d7c8..d96d91e59 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandler.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandler.java @@ -44,7 +44,7 @@ public class JVMMetricsServiceHandler extends JVMMetricsServiceGrpc.JVMMetricsSe StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); request.getMetricsList().forEach(metric -> { long time = TimeBucketUtils.INSTANCE.getSecondTimeBucket(metric.getTime()); - senToInstanceHeartBeatPersistenceWorker(context, applicationInstanceId, time); + senToInstanceHeartBeatPersistenceWorker(context, applicationInstanceId, metric.getTime()); sendToCpuMetricPersistenceWorker(context, applicationInstanceId, time, metric.getCpu()); sendToMemoryMetricPersistenceWorker(context, applicationInstanceId, time, metric.getMemoryList()); sendToMemoryPoolMetricPersistenceWorker(context, applicationInstanceId, time, metric.getMemoryPoolList()); @@ -59,8 +59,8 @@ public class JVMMetricsServiceHandler extends JVMMetricsServiceGrpc.JVMMetricsSe long heartBeatTime) { InstanceHeartBeatDataDefine.InstanceHeartBeat heartBeat = new InstanceHeartBeatDataDefine.InstanceHeartBeat(); heartBeat.setId(String.valueOf(applicationInstanceId)); - heartBeat.setHeartbeatTime(heartBeatTime); - heartBeat.setApplicationInstanceId(applicationInstanceId); + heartBeat.setHeartBeatTime(heartBeatTime); + heartBeat.setInstanceId(applicationInstanceId); try { logger.debug("send to instance heart beat persistence worker, id: {}", heartBeat.getId()); context.getClusterWorkerContext().lookup(InstHeartBeatPersistenceWorker.WorkerRole.INSTANCE).tell(heartBeat.toData()); diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/InstHeartBeatPersistenceWorker.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/InstHeartBeatPersistenceWorker.java index 289b13326..02ae71c7e 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/InstHeartBeatPersistenceWorker.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/InstHeartBeatPersistenceWorker.java @@ -1,6 +1,6 @@ package org.skywalking.apm.collector.agentjvm.worker.heartbeat; -import org.skywalking.apm.collector.agentjvm.worker.heartbeat.dao.InstanceHeartBeatEsDAO; +import org.skywalking.apm.collector.agentjvm.worker.heartbeat.dao.IInstanceHeartBeatDAO; import org.skywalking.apm.collector.agentjvm.worker.heartbeat.define.InstanceHeartBeatDataDefine; import org.skywalking.apm.collector.storage.dao.DAOContainer; import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorkerProvider; @@ -27,11 +27,11 @@ public class InstHeartBeatPersistenceWorker extends PersistenceWorker { } @Override protected boolean needMergeDBData() { - return false; + return true; } @Override protected IPersistenceDAO persistenceDAO() { - return (IPersistenceDAO)DAOContainer.INSTANCE.get(InstanceHeartBeatEsDAO.class.getName()); + return (IPersistenceDAO)DAOContainer.INSTANCE.get(IInstanceHeartBeatDAO.class.getName()); } public static class Factory extends AbstractLocalAsyncWorkerProvider { diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/dao/InstanceHeartBeatEsDAO.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/dao/InstanceHeartBeatEsDAO.java index 3ced6633a..4a011c2d7 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/dao/InstanceHeartBeatEsDAO.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/dao/InstanceHeartBeatEsDAO.java @@ -2,30 +2,47 @@ package org.skywalking.apm.collector.agentjvm.worker.heartbeat.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.core.framework.UnexpectedException; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; import org.skywalking.apm.collector.storage.table.register.InstanceTable; 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; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author pengys5 */ public class InstanceHeartBeatEsDAO extends EsDAO implements IInstanceHeartBeatDAO, IPersistenceDAO { + private final Logger logger = LoggerFactory.getLogger(InstanceHeartBeatEsDAO.class); + @Override public Data get(String id, DataDefine dataDefine) { - return null; + GetResponse getResponse = getClient().prepareGet(InstanceTable.TABLE, id).get(); + if (getResponse.isExists()) { + Data data = dataDefine.build(id); + Map source = getResponse.getSource(); + data.setDataInteger(0, (Integer)source.get(InstanceTable.COLUMN_INSTANCE_ID)); + data.setDataLong(0, (Long)source.get(InstanceTable.COLUMN_HEARTBEAT_TIME)); + logger.debug("id: {} is exists", id); + return data; + } else { + logger.debug("id: {} is not exists", id); + return null; + } } @Override public IndexRequestBuilder prepareBatchInsert(Data data) { - return null; + throw new UnexpectedException("There is no need to merge stream data with database data."); } @Override public UpdateRequestBuilder prepareBatchUpdate(Data data) { Map source = new HashMap<>(); - source.put(InstanceTable.COLUMN_REGISTER_TIME, data.getDataLong(0)); + source.put(InstanceTable.COLUMN_HEARTBEAT_TIME, data.getDataLong(0)); return getClient().prepareUpdate(InstanceTable.TABLE, data.getDataString(0)).setDoc(source); } } diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/define/InstanceHeartBeatDataDefine.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/define/InstanceHeartBeatDataDefine.java index b91275fa4..7af4f5b9e 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/define/InstanceHeartBeatDataDefine.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/define/InstanceHeartBeatDataDefine.java @@ -1,6 +1,5 @@ package org.skywalking.apm.collector.agentjvm.worker.heartbeat.define; -import org.skywalking.apm.collector.core.framework.UnexpectedException; 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; @@ -22,27 +21,35 @@ public class InstanceHeartBeatDataDefine extends DataDefine { @Override protected void attributeDefine() { addAttribute(0, new Attribute(InstanceTable.COLUMN_ID, AttributeType.STRING, new NonOperation())); - addAttribute(1, new Attribute(InstanceTable.COLUMN_INSTANCE_ID, AttributeType.INTEGER, new NonOperation())); + addAttribute(1, new Attribute(InstanceTable.COLUMN_INSTANCE_ID, AttributeType.INTEGER, new CoverOperation())); addAttribute(2, new Attribute(InstanceTable.COLUMN_HEARTBEAT_TIME, AttributeType.LONG, new CoverOperation())); } @Override public Object deserialize(RemoteData remoteData) { - throw new UnexpectedException("instance heart beat data did not need send to remote worker."); + String id = remoteData.getDataStrings(0); + int instanceId = remoteData.getDataIntegers(0); + long heartBeatTime = remoteData.getDataLongs(0); + return new InstanceHeartBeat(id, heartBeatTime, instanceId); } @Override public RemoteData serialize(Object object) { - throw new UnexpectedException("instance heart beat data did not need send to remote worker."); + InstanceHeartBeat instanceHeartBeat = (InstanceHeartBeat)object; + RemoteData.Builder builder = RemoteData.newBuilder(); + builder.addDataStrings(instanceHeartBeat.getId()); + builder.addDataIntegers(instanceHeartBeat.getInstanceId()); + builder.addDataLongs(instanceHeartBeat.getHeartBeatTime()); + return builder.build(); } public static class InstanceHeartBeat implements Transform { private String id; - private int applicationInstanceId; - private long heartbeatTime; + private long heartBeatTime; + private int instanceId; - public InstanceHeartBeat(String id, int applicationInstanceId, long heartbeatTime) { + public InstanceHeartBeat(String id, long heartBeatTime, int instanceId) { this.id = id; - this.applicationInstanceId = applicationInstanceId; - this.heartbeatTime = heartbeatTime; + this.heartBeatTime = heartBeatTime; + this.instanceId = instanceId; } public InstanceHeartBeat() { @@ -52,40 +59,40 @@ public class InstanceHeartBeatDataDefine extends DataDefine { InstanceHeartBeatDataDefine define = new InstanceHeartBeatDataDefine(); Data data = define.build(id); data.setDataString(0, this.id); - data.setDataInteger(0, this.applicationInstanceId); - data.setDataLong(0, this.heartbeatTime); + data.setDataInteger(0, this.instanceId); + data.setDataLong(0, this.heartBeatTime); return data; } @Override public InstanceHeartBeat toSelf(Data data) { this.id = data.getDataString(0); - this.applicationInstanceId = data.getDataInteger(0); - this.heartbeatTime = data.getDataLong(0); + this.instanceId = data.getDataInteger(0); + this.heartBeatTime = data.getDataLong(0); return this; } - public void setId(String id) { - this.id = id; - } - - public void setApplicationInstanceId(int applicationInstanceId) { - this.applicationInstanceId = applicationInstanceId; - } - public String getId() { return id; } - public int getApplicationInstanceId() { - return applicationInstanceId; + public void setId(String id) { + this.id = id; } - public long getHeartbeatTime() { - return heartbeatTime; + public long getHeartBeatTime() { + return heartBeatTime; } - public void setHeartbeatTime(long heartbeatTime) { - this.heartbeatTime = heartbeatTime; + public void setHeartBeatTime(long heartBeatTime) { + this.heartBeatTime = heartBeatTime; + } + + public int getInstanceId() { + return instanceId; + } + + public void setInstanceId(int instanceId) { + this.instanceId = instanceId; } } } diff --git a/apm-collector/apm-collector-agentjvm/src/main/resources/META-INF/defines/es_dao.define b/apm-collector/apm-collector-agentjvm/src/main/resources/META-INF/defines/es_dao.define index 7b18e4a00..13e6077c0 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/resources/META-INF/defines/es_dao.define +++ b/apm-collector/apm-collector-agentjvm/src/main/resources/META-INF/defines/es_dao.define @@ -1,4 +1,5 @@ org.skywalking.apm.collector.agentjvm.worker.cpu.dao.CpuMetricEsDAO org.skywalking.apm.collector.agentjvm.worker.memory.dao.MemoryMetricEsDAO org.skywalking.apm.collector.agentjvm.worker.memorypool.dao.MemoryPoolMetricEsDAO -org.skywalking.apm.collector.agentjvm.worker.gc.dao.GCMetricEsDAO \ No newline at end of file +org.skywalking.apm.collector.agentjvm.worker.gc.dao.GCMetricEsDAO +org.skywalking.apm.collector.agentjvm.worker.heartbeat.dao.InstanceHeartBeatEsDAO \ No newline at end of file diff --git a/apm-collector/apm-collector-agentjvm/src/main/resources/META-INF/defines/h2_dao.define b/apm-collector/apm-collector-agentjvm/src/main/resources/META-INF/defines/h2_dao.define index 2c5bdb85e..c6eb797fb 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/resources/META-INF/defines/h2_dao.define +++ b/apm-collector/apm-collector-agentjvm/src/main/resources/META-INF/defines/h2_dao.define @@ -1,4 +1,5 @@ org.skywalking.apm.collector.agentjvm.worker.cpu.dao.CpuMetricH2DAO org.skywalking.apm.collector.agentjvm.worker.memory.dao.MemoryMetricH2DAO org.skywalking.apm.collector.agentjvm.worker.memorypool.dao.MemoryPoolMetricH2DAO -org.skywalking.apm.collector.agentjvm.worker.gc.dao.GCMetricH2DAO \ No newline at end of file +org.skywalking.apm.collector.agentjvm.worker.gc.dao.GCMetricH2DAO +org.skywalking.apm.collector.agentjvm.worker.heartbeat.dao.InstanceHeartBeatH2DAO \ No newline at end of file diff --git a/apm-collector/apm-collector-agentjvm/src/test/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandlerTestCase.java b/apm-collector/apm-collector-agentjvm/src/test/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandlerTestCase.java index 849dc363a..3b2121e85 100644 --- a/apm-collector/apm-collector-agentjvm/src/test/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandlerTestCase.java +++ b/apm-collector/apm-collector-agentjvm/src/test/java/org/skywalking/apm/collector/agentjvm/grpc/handler/JVMMetricsServiceHandlerTestCase.java @@ -21,14 +21,25 @@ public class JVMMetricsServiceHandlerTestCase { private final Logger logger = LoggerFactory.getLogger(JVMMetricsServiceHandlerTestCase.class); - private JVMMetricsServiceGrpc.JVMMetricsServiceBlockingStub stub; + private static JVMMetricsServiceGrpc.JVMMetricsServiceBlockingStub stub; - public void test() { + public static void main(String[] args) { ManagedChannel channel = ManagedChannelBuilder.forAddress("localhost", 11800).usePlaintext(true).build(); stub = JVMMetricsServiceGrpc.newBlockingStub(channel); + buildJvmMetric(2); + buildJvmMetric(3); + + try { + Thread.sleep(2000); + } catch (InterruptedException e) { + e.printStackTrace(); + } + } + + public static void buildJvmMetric(int instanceId) { JVMMetrics.Builder jvmMetricsBuilder = JVMMetrics.newBuilder(); - jvmMetricsBuilder.setApplicationInstanceId(1); + jvmMetricsBuilder.setApplicationInstanceId(instanceId); JVMMetric.Builder jvmMetric = JVMMetric.newBuilder(); jvmMetric.setTime(System.currentTimeMillis()); @@ -41,13 +52,13 @@ public class JVMMetricsServiceHandlerTestCase { stub.collect(jvmMetricsBuilder.build()); } - private void buildCpuMetric(JVMMetric.Builder jvmMetric) { + private static void buildCpuMetric(JVMMetric.Builder jvmMetric) { CPU.Builder cpuBuilder = CPU.newBuilder(); cpuBuilder.setUsagePercent(70); jvmMetric.setCpu(cpuBuilder); } - private void buildMemoryMetric(JVMMetric.Builder jvmMetric) { + private static void buildMemoryMetric(JVMMetric.Builder jvmMetric) { Memory.Builder builder_1 = Memory.newBuilder(); builder_1.setIsHeap(true); builder_1.setInit(20); @@ -65,7 +76,7 @@ public class JVMMetricsServiceHandlerTestCase { jvmMetric.addMemory(builder_2.build()); } - private void buildMemoryPoolMetric(JVMMetric.Builder jvmMetric) { + private static void buildMemoryPoolMetric(JVMMetric.Builder jvmMetric) { MemoryPool.Builder builder_1 = MemoryPool.newBuilder(); builder_1.setType(PoolType.NEWGEN_USAGE); builder_1.setIsHeap(true); @@ -76,7 +87,7 @@ public class JVMMetricsServiceHandlerTestCase { jvmMetric.addMemoryPool(builder_1.build()); } - private void buildGcMetric(JVMMetric.Builder jvmMetric) { + private static void buildGcMetric(JVMMetric.Builder jvmMetric) { GC.Builder gcBuilder = GC.newBuilder(); gcBuilder.setPhrase(GCPhrase.NEW); gcBuilder.setCount(2); diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java index 9929f22bc..9f48383ed 100644 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java @@ -26,7 +26,7 @@ public class InstanceIDService { if (instanceId == 0) { StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); - InstanceDataDefine.Instance instance = new InstanceDataDefine.Instance("0", applicationId, agentUUID, registerTime, 0); + InstanceDataDefine.Instance instance = new InstanceDataDefine.Instance("0", applicationId, agentUUID, registerTime, 0, registerTime); try { context.getClusterWorkerContext().lookup(ApplicationRegisterRemoteWorker.WorkerRole.INSTANCE).tell(instance); } catch (WorkerNotFoundException | WorkerInvokeException e) { @@ -46,7 +46,7 @@ public class InstanceIDService { logger.debug("instance recover, instance id: {}, application id: {}, register time: {}", instanceId, applicationId, registerTime); IInstanceDAO dao = (IInstanceDAO)DAOContainer.INSTANCE.get(IInstanceDAO.class.getName()); - InstanceDataDefine.Instance instance = new InstanceDataDefine.Instance(String.valueOf(instanceId), applicationId, "", registerTime, instanceId); + InstanceDataDefine.Instance instance = new InstanceDataDefine.Instance(String.valueOf(instanceId), applicationId, "", registerTime, instanceId, registerTime); dao.save(instance); } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceEsTableDefine.java new file mode 100644 index 000000000..c571b38db --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceEsTableDefine.java @@ -0,0 +1,35 @@ +package org.skywalking.apm.collector.agentstream.worker.instance.performance.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.instance.InstPerformanceTable; + +/** + * @author pengys5 + */ +public class InstPerformanceEsTableDefine extends ElasticSearchTableDefine { + + public InstPerformanceEsTableDefine() { + super(InstPerformanceTable.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(InstPerformanceTable.COLUMN_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(InstPerformanceTable.COLUMN_INSTANCE_ID, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(InstPerformanceTable.COLUMN_CALL_TIMES, ElasticSearchColumnDefine.Type.Integer.name())); + addColumn(new ElasticSearchColumnDefine(InstPerformanceTable.COLUMN_COST_TOTAL, ElasticSearchColumnDefine.Type.Long.name())); + addColumn(new ElasticSearchColumnDefine(InstPerformanceTable.COLUMN_TIME_BUCKET, ElasticSearchColumnDefine.Type.Long.name())); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceH2TableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceH2TableDefine.java new file mode 100644 index 000000000..b01c8e784 --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceH2TableDefine.java @@ -0,0 +1,24 @@ +package org.skywalking.apm.collector.agentstream.worker.instance.performance.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.instance.InstPerformanceTable; + +/** + * @author pengys5 + */ +public class InstPerformanceH2TableDefine extends H2TableDefine { + + public InstPerformanceH2TableDefine() { + super(InstPerformanceTable.TABLE); + } + + @Override public void initialize() { + addColumn(new H2ColumnDefine(InstPerformanceTable.COLUMN_ID, H2ColumnDefine.Type.Varchar.name())); + addColumn(new H2ColumnDefine(InstPerformanceTable.COLUMN_APPLICATION_ID, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(InstPerformanceTable.COLUMN_INSTANCE_ID, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(InstPerformanceTable.COLUMN_CALL_TIMES, H2ColumnDefine.Type.Int.name())); + addColumn(new H2ColumnDefine(InstPerformanceTable.COLUMN_COST_TOTAL, H2ColumnDefine.Type.Bigint.name())); + addColumn(new H2ColumnDefine(InstPerformanceTable.COLUMN_TIME_BUCKET, H2ColumnDefine.Type.Bigint.name())); + } +} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationEsTableDefine.java index d762a6346..3ac13b5a9 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationEsTableDefine.java @@ -14,7 +14,7 @@ public class ApplicationEsTableDefine extends ElasticSearchTableDefine { } @Override public int refreshInterval() { - return 0; + return 2; } @Override public int numberOfShards() { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceDataDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceDataDefine.java index 0d2f5a93a..be9a0bc8d 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceDataDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceDataDefine.java @@ -32,7 +32,8 @@ public class InstanceDataDefine extends DataDefine { String agentUUID = remoteData.getDataStrings(1); int instanceId = remoteData.getDataIntegers(1); long registerTime = remoteData.getDataLongs(0); - return new Instance(id, applicationId, agentUUID, registerTime, instanceId); + long heartBeatTime = remoteData.getDataLongs(1); + return new Instance(id, applicationId, agentUUID, registerTime, instanceId, heartBeatTime); } @Override public RemoteData serialize(Object object) { @@ -42,6 +43,7 @@ public class InstanceDataDefine extends DataDefine { builder.addDataIntegers(instance.getApplicationId()); builder.addDataStrings(instance.getAgentUUID()); builder.addDataLongs(instance.getRegisterTime()); + builder.addDataLongs(instance.getHeartBeatTime()); return builder.build(); } @@ -51,13 +53,16 @@ public class InstanceDataDefine extends DataDefine { private String agentUUID; private long registerTime; private int instanceId; + private long heartBeatTime; - public Instance(String id, int applicationId, String agentUUID, long registerTime, int instanceId) { + public Instance(String id, int applicationId, String agentUUID, long registerTime, int instanceId, + long heartBeatTime) { this.id = id; this.applicationId = applicationId; this.agentUUID = agentUUID; this.registerTime = registerTime; this.instanceId = instanceId; + this.heartBeatTime = heartBeatTime; } public String getId() { @@ -99,5 +104,13 @@ public class InstanceDataDefine extends DataDefine { public void setInstanceId(int instanceId) { this.instanceId = instanceId; } + + public long getHeartBeatTime() { + return heartBeatTime; + } + + public void setHeartBeatTime(long heartBeatTime) { + this.heartBeatTime = heartBeatTime; + } } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceEsTableDefine.java index 0635869d6..3d0ba91d1 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceEsTableDefine.java @@ -14,7 +14,7 @@ public class InstanceEsTableDefine extends ElasticSearchTableDefine { } @Override public int refreshInterval() { - return 0; + return 2; } @Override public int numberOfShards() { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceEsDAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceEsDAO.java index 69456c7f1..31ccdacbe 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceEsDAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceEsDAO.java @@ -13,9 +13,9 @@ 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.storage.table.register.InstanceTable; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.storage.table.register.InstanceTable; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -62,6 +62,7 @@ public class InstanceEsDAO extends EsDAO implements IInstanceDAO { source.put(InstanceTable.COLUMN_APPLICATION_ID, instance.getApplicationId()); source.put(InstanceTable.COLUMN_AGENT_UUID, instance.getAgentUUID()); source.put(InstanceTable.COLUMN_REGISTER_TIME, instance.getRegisterTime()); + source.put(InstanceTable.COLUMN_HEARTBEAT_TIME, instance.getHeartBeatTime()); IndexResponse response = client.prepareIndex(InstanceTable.TABLE, instance.getId()).setSource(source).setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE).get(); logger.debug("save instance register info, application id: {}, agentUUID: {}, status: {}", instance.getApplicationId(), instance.getAgentUUID(), response.status().name()); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameEsTableDefine.java index c1a90efdf..455ab4375 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameEsTableDefine.java @@ -14,7 +14,7 @@ public class ServiceNameEsTableDefine extends ElasticSearchTableDefine { } @Override public int refreshInterval() { - return 0; + return 2; } @Override public int numberOfShards() { diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define index c462b86ec..6caafb6b7 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define @@ -32,4 +32,7 @@ org.skywalking.apm.collector.agentstream.worker.service.entry.define.ServiceEntr org.skywalking.apm.collector.agentstream.worker.service.entry.define.ServiceEntryH2TableDefine org.skywalking.apm.collector.agentstream.worker.serviceref.reference.define.ServiceRefEsTableDefine -org.skywalking.apm.collector.agentstream.worker.serviceref.reference.define.ServiceRefH2TableDefine \ No newline at end of file +org.skywalking.apm.collector.agentstream.worker.serviceref.reference.define.ServiceRefH2TableDefine + +org.skywalking.apm.collector.agentstream.worker.instance.performance.define.InstPerformanceEsTableDefine +org.skywalking.apm.collector.agentstream.worker.instance.performance.define.InstPerformanceH2TableDefine \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java index c153190ec..2ffcad2cf 100644 --- a/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java +++ b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java @@ -22,9 +22,9 @@ public class SegmentPost { InstanceEsDAO instanceEsDAO = new InstanceEsDAO(); instanceEsDAO.setClient(client); - InstanceDataDefine.Instance consumerInstance = new InstanceDataDefine.Instance("2", 2, "dubbox-consumer", 1501858094526L, 2); + InstanceDataDefine.Instance consumerInstance = new InstanceDataDefine.Instance("2", 2, "dubbox-consumer", 1501858094526L, 2, 1501858094526L); instanceEsDAO.save(consumerInstance); - InstanceDataDefine.Instance providerInstance = new InstanceDataDefine.Instance("3", 3, "dubbox-provider", 1501858094526L, 3); + InstanceDataDefine.Instance providerInstance = new InstanceDataDefine.Instance("3", 3, "dubbox-provider", 1501858094526L, 3, 1501858094526L); instanceEsDAO.save(providerInstance); ApplicationEsDAO applicationEsDAO = new ApplicationEsDAO(); diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceEsDAO.java index 82ae92b7d..22ee0315a 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceEsDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceEsDAO.java @@ -38,7 +38,7 @@ public class InstPerformanceEsDAO extends EsDAO implements IInstPerformanceDAO { for (SearchHit searchHit : searchHits) { int instanceId = (Integer)searchHit.getSource().get(InstPerformanceTable.COLUMN_INSTANCE_ID); int callTimes = (Integer)searchHit.getSource().get(InstPerformanceTable.COLUMN_CALL_TIMES); - long costTotal = (Long)searchHit.getSource().get(InstPerformanceTable.COLUMN_COST_TOTAL); + long costTotal = ((Number)searchHit.getSource().get(InstPerformanceTable.COLUMN_COST_TOTAL)).longValue(); instPerformances.add(new InstPerformance(instanceId, callTimes, costTotal)); } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceEsDAO.java index 88f4103f3..ecef9bf0e 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceEsDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceEsDAO.java @@ -65,13 +65,19 @@ public class InstanceEsDAO extends EsDAO implements IInstanceDAO { } @Override public List getApplications(long time) { + logger.debug("application list get, time: {}", time); SearchRequestBuilder searchRequestBuilder = getClient().prepareSearch(InstanceTable.TABLE); searchRequestBuilder.setTypes(InstanceTable.TABLE_TYPE); searchRequestBuilder.setSearchType(SearchType.DFS_QUERY_THEN_FETCH); - RangeQueryBuilder rangeQueryBuilder = QueryBuilders.rangeQuery(InstanceTable.COLUMN_HEARTBEAT_TIME).gt(time); + BoolQueryBuilder boolQueryBuilder = QueryBuilders.boolQuery(); + RangeQueryBuilder heartBeatRangeQueryBuilder = QueryBuilders.rangeQuery(InstanceTable.COLUMN_HEARTBEAT_TIME).gte(time); + RangeQueryBuilder registerRangeQueryBuilder = QueryBuilders.rangeQuery(InstanceTable.COLUMN_REGISTER_TIME).lte(time); - searchRequestBuilder.setQuery(rangeQueryBuilder); + boolQueryBuilder.must().add(registerRangeQueryBuilder); + boolQueryBuilder.must().add(heartBeatRangeQueryBuilder); + + searchRequestBuilder.setQuery(boolQueryBuilder); searchRequestBuilder.setSize(0); searchRequestBuilder.addAggregation(AggregationBuilders.terms(InstanceTable.COLUMN_APPLICATION_ID).field(InstanceTable.COLUMN_APPLICATION_ID).size(100)); @@ -81,6 +87,7 @@ public class InstanceEsDAO extends EsDAO implements IInstanceDAO { List applications = new LinkedList<>(); for (Terms.Bucket entry : genders.getBuckets()) { Integer applicationId = entry.getKeyAsNumber().intValue(); + logger.debug("applicationId: {}", applicationId); long instanceCount = entry.getDocCount(); applications.add(new Application(applicationId, instanceCount)); } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceH2DAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceH2DAO.java index b5568f662..2e1badda7 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceH2DAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstanceH2DAO.java @@ -1,5 +1,6 @@ package org.skywalking.apm.collector.ui.dao; +import java.util.List; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; /** @@ -13,4 +14,8 @@ public class InstanceH2DAO extends H2DAO implements IInstanceDAO { @Override public Long instanceLastHeartBeatTime(long applicationInstanceId) { return null; } + + @Override public List getApplications(long time) { + return null; + } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/UIJettyModuleDefine.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/UIJettyModuleDefine.java index 9d0e73bb3..cfad93a70 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/UIJettyModuleDefine.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/UIJettyModuleDefine.java @@ -15,6 +15,10 @@ import org.skywalking.apm.collector.ui.jetty.handler.SpanGetHandler; import org.skywalking.apm.collector.ui.jetty.handler.TraceDagGetHandler; import org.skywalking.apm.collector.ui.jetty.handler.TraceStackGetHandler; import org.skywalking.apm.collector.ui.jetty.handler.UIJettyServerHandler; +import org.skywalking.apm.collector.ui.jetty.handler.instancehealth.ApplicationsGetHandler; +import org.skywalking.apm.collector.ui.jetty.handler.instancehealth.InstanceHealthGetHandler; +import org.skywalking.apm.collector.ui.jetty.handler.time.AllInstanceLastTimeGetHandler; +import org.skywalking.apm.collector.ui.jetty.handler.time.InstanceLastTimeGetHandler; /** * @author pengys5 @@ -54,6 +58,10 @@ public class UIJettyModuleDefine extends UIModuleDefine { handlers.add(new SegmentTopGetHandler()); handlers.add(new TraceStackGetHandler()); handlers.add(new SpanGetHandler()); + handlers.add(new InstanceLastTimeGetHandler()); + handlers.add(new AllInstanceLastTimeGetHandler()); + handlers.add(new ApplicationsGetHandler()); + handlers.add(new InstanceHealthGetHandler()); return handlers; } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceHealthService.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceHealthService.java index 72ea89207..66ac5c3bd 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceHealthService.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceHealthService.java @@ -5,7 +5,6 @@ import com.google.gson.JsonObject; import java.util.List; import org.skywalking.apm.collector.storage.dao.DAOContainer; import org.skywalking.apm.collector.ui.cache.ApplicationCache; -import org.skywalking.apm.collector.ui.dao.GCMetricEsDAO; import org.skywalking.apm.collector.ui.dao.IGCMetricDAO; import org.skywalking.apm.collector.ui.dao.IInstPerformanceDAO; import org.skywalking.apm.collector.ui.dao.IInstanceDAO; @@ -35,6 +34,7 @@ public class InstanceHealthService { applicationJson.addProperty("applicationId", application.getApplicationId()); applicationJson.addProperty("applicationCode", applicationCode); applicationJson.addProperty("instanceCount", application.getCount()); + applicationArray.add(applicationJson); }); return response; @@ -53,7 +53,7 @@ public class InstanceHealthService { response.addProperty("applicationId", applicationId); response.add("appInstances", instances); - GCMetricEsDAO gcMetricEsDAO = (GCMetricEsDAO)DAOContainer.INSTANCE.get(GCMetricEsDAO.class.getName()); + IGCMetricDAO gcMetricDAO = (IGCMetricDAO)DAOContainer.INSTANCE.get(IGCMetricDAO.class.getName()); performances.forEach(instance -> { JsonObject instanceJson = new JsonObject(); instanceJson.addProperty("id", instance.getInstanceId()); @@ -74,7 +74,7 @@ public class InstanceHealthService { instanceJson.addProperty("status", 0); - IGCMetricDAO.GCCount gcCount = gcMetricEsDAO.getGCCount(timestamp, instance.getInstanceId()); + IGCMetricDAO.GCCount gcCount = gcMetricDAO.getGCCount(timestamp, instance.getInstanceId()); instanceJson.addProperty("ygc", gcCount.getYoung()); instanceJson.addProperty("ogc", gcCount.getOld());