From bc38a56da0a00dc309e4b48d1c1fa8071c940ce4 Mon Sep 17 00:00:00 2001 From: clevertension Date: Thu, 12 Oct 2017 21:19:53 +0800 Subject: [PATCH] add h2 support ga --- .../worker/cpu/dao/CpuMetricH2DAO.java | 15 +- .../agentjvm/worker/gc/dao/GCMetricH2DAO.java | 15 +- .../heartbeat/dao/InstanceHeartBeatH2DAO.java | 45 +++--- .../worker/memory/dao/MemoryMetricH2DAO.java | 14 +- .../memorypool/dao/MemoryPoolMetricH2DAO.java | 14 +- .../application/dao/ApplicationH2DAO.java | 11 +- .../worker/instance/dao/InstanceH2DAO.java | 72 +++++++++- .../worker/global/dao/GlobalTraceH2DAO.java | 31 +++- .../performance/dao/InstPerformanceH2DAO.java | 61 +++++++- .../component/dao/NodeComponentH2DAO.java | 70 ++++++++- .../node/mapping/dao/NodeMappingH2DAO.java | 57 +++++++- .../noderef/dao/NodeReferenceH2DAO.java | 10 +- .../segment/cost/dao/SegmentCostH2DAO.java | 34 +++-- .../segment/origin/dao/SegmentH2DAO.java | 21 ++- .../service/entry/dao/ServiceEntryH2DAO.java | 35 +++-- .../serviceref/dao/ServiceReferenceH2DAO.java | 78 +++++----- .../src/main/resources/application.yml | 8 +- .../apm/collector/client/h2/H2Client.java | 28 ++-- .../collector/storage/h2/dao/BatchH2DAO.java | 63 +++++++- .../apm/collector/storage/h2/dao/H2DAO.java | 34 ++++- .../storage/h2/define/H2SqlEntity.java | 22 +++ .../storage/h2/define/H2StorageInstaller.java | 9 +- .../collector/ui/dao/ApplicationH2DAO.java | 25 +++- .../apm/collector/ui/dao/CpuMetricH2DAO.java | 81 +++++++++++ .../apm/collector/ui/dao/GCMetricH2DAO.java | 9 +- .../ui/dao/InstPerformanceH2DAO.java | 132 ++++++++++++++++- .../apm/collector/ui/dao/InstanceH2DAO.java | 134 ++++++++++++++++-- .../collector/ui/dao/NodeComponentH2DAO.java | 52 ++++++- .../collector/ui/dao/NodeMappingH2DAO.java | 48 ++++++- .../collector/ui/dao/NodeReferenceH2DAO.java | 63 +++++++- .../ui/service/InstanceHealthService.java | 8 +- 31 files changed, 1093 insertions(+), 206 deletions(-) create mode 100644 apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2SqlEntity.java create mode 100644 apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/CpuMetricH2DAO.java diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/cpu/dao/CpuMetricH2DAO.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/cpu/dao/CpuMetricH2DAO.java index 065f4c918..169dec800 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/cpu/dao/CpuMetricH2DAO.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/cpu/dao/CpuMetricH2DAO.java @@ -3,8 +3,8 @@ package org.skywalking.apm.collector.agentjvm.worker.cpu.dao; import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.storage.define.jvm.CpuMetricTable; -import org.skywalking.apm.collector.storage.define.jvm.GCMetricTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.storage.h2.define.H2SqlEntity; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -15,23 +15,28 @@ import java.util.Map; /** * @author pengys5 */ -public class CpuMetricH2DAO extends H2DAO implements ICpuMetricDAO, IPersistenceDAO, Map> { +public class CpuMetricH2DAO extends H2DAO implements ICpuMetricDAO, IPersistenceDAO { private final Logger logger = LoggerFactory.getLogger(CpuMetricH2DAO.class); @Override public Data get(String id, DataDefine dataDefine) { return null; } - @Override public Map prepareBatchInsert(Data data) { + @Override public H2SqlEntity prepareBatchInsert(Data data) { + H2SqlEntity entity = new H2SqlEntity(); Map source = new HashMap<>(); + source.put("id", data.getDataString(0)); source.put(CpuMetricTable.COLUMN_INSTANCE_ID, data.getDataInteger(0)); source.put(CpuMetricTable.COLUMN_USAGE_PERCENT, data.getDataDouble(0)); source.put(CpuMetricTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); logger.debug("prepare cpu metric batch insert, id: {}", data.getDataString(0)); - return source; + String sql = getBatchInsertSql(CpuMetricTable.TABLE, source.keySet()); + entity.setSql(sql); + entity.setParams(source.values().toArray(new Object[0])); + return entity; } - @Override public Map prepareBatchUpdate(Data data) { + @Override public H2SqlEntity prepareBatchUpdate(Data data) { return null; } } diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/dao/GCMetricH2DAO.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/dao/GCMetricH2DAO.java index b58a40158..268762ec9 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/dao/GCMetricH2DAO.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/dao/GCMetricH2DAO.java @@ -3,8 +3,8 @@ package org.skywalking.apm.collector.agentjvm.worker.gc.dao; import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.storage.define.jvm.GCMetricTable; -import org.skywalking.apm.collector.storage.define.jvm.MemoryMetricTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.storage.h2.define.H2SqlEntity; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import java.util.HashMap; @@ -13,23 +13,28 @@ import java.util.Map; /** * @author pengys5 */ -public class GCMetricH2DAO extends H2DAO implements IGCMetricDAO, IPersistenceDAO, Map> { +public class GCMetricH2DAO extends H2DAO implements IGCMetricDAO, IPersistenceDAO { @Override public Data get(String id, DataDefine dataDefine) { return null; } - @Override public Map prepareBatchInsert(Data data) { + @Override public H2SqlEntity prepareBatchInsert(Data data) { + H2SqlEntity entity = new H2SqlEntity(); Map source = new HashMap<>(); + source.put("id", data.getDataString(0)); source.put(GCMetricTable.COLUMN_INSTANCE_ID, data.getDataInteger(0)); source.put(GCMetricTable.COLUMN_PHRASE, data.getDataInteger(1)); source.put(GCMetricTable.COLUMN_COUNT, data.getDataLong(0)); source.put(GCMetricTable.COLUMN_TIME, data.getDataLong(1)); source.put(GCMetricTable.COLUMN_TIME_BUCKET, data.getDataLong(2)); - return source; + String sql = getBatchInsertSql(GCMetricTable.TABLE, source.keySet()); + entity.setSql(sql); + entity.setParams(source.values().toArray(new Object[0])); + return entity; } - @Override public Map prepareBatchUpdate(Data data) { + @Override public H2SqlEntity prepareBatchUpdate(Data data) { return null; } } diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/dao/InstanceHeartBeatH2DAO.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/dao/InstanceHeartBeatH2DAO.java index fc10e863e..9a6cd29ca 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/dao/InstanceHeartBeatH2DAO.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/heartbeat/dao/InstanceHeartBeatH2DAO.java @@ -1,8 +1,5 @@ package org.skywalking.apm.collector.agentjvm.worker.heartbeat.dao; -import org.elasticsearch.action.get.GetResponse; -import org.elasticsearch.action.index.IndexRequestBuilder; -import org.elasticsearch.action.update.UpdateRequestBuilder; import org.skywalking.apm.collector.client.h2.H2Client; import org.skywalking.apm.collector.client.h2.H2ClientException; import org.skywalking.apm.collector.core.framework.UnexpectedException; @@ -10,48 +7,52 @@ import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.storage.define.register.InstanceTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.storage.h2.define.H2SqlEntity; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.sql.ResultSet; import java.sql.SQLException; -import java.util.HashMap; -import java.util.Map; +import java.text.MessageFormat; +import java.util.*; /** * @author pengys5 */ -public class InstanceHeartBeatH2DAO extends H2DAO implements IInstanceHeartBeatDAO, IPersistenceDAO, Map> { +public class InstanceHeartBeatH2DAO extends H2DAO implements IInstanceHeartBeatDAO, IPersistenceDAO { private final Logger logger = LoggerFactory.getLogger(InstanceHeartBeatEsDAO.class); - + private static final String GET_INSTANCE_HEARTBEAT_SQL = "select * from {0} where {1} = ?"; @Override public Data get(String id, DataDefine dataDefine) { H2Client client = getClient(); - String sql = "select " + InstanceTable.COLUMN_INSTANCE_ID + "," + InstanceTable.COLUMN_HEARTBEAT_TIME + - " from " + InstanceTable.TABLE + " where " + InstanceTable.COLUMN_INSTANCE_ID + "=?"; - Object[] params = new Object[] {id}; - ResultSet rs = null; - try { - rs = client.executeQuery(sql, params); - Data data = dataDefine.build(id); - data.setDataInteger(0, rs.getInt(1)); - data.setDataLong(0, rs.getLong(2)); - return data; + String sql = MessageFormat.format(GET_INSTANCE_HEARTBEAT_SQL, InstanceTable.TABLE, InstanceTable.COLUMN_INSTANCE_ID); + Object[] params = new Object[]{id}; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + Data data = dataDefine.build(id); + data.setDataInteger(0, rs.getInt(InstanceTable.COLUMN_INSTANCE_ID)); + data.setDataLong(0, rs.getLong(InstanceTable.COLUMN_HEARTBEAT_TIME)); + return data; + } } catch (SQLException | H2ClientException e) { logger.error(e.getMessage(), e); - } finally { - client.closeResultSet(rs); } return null; } - @Override public Map prepareBatchInsert(Data data) { + @Override public H2SqlEntity prepareBatchInsert(Data data) { throw new UnexpectedException("There is no need to merge stream data with database data."); } - @Override public Map prepareBatchUpdate(Data data) { + @Override public H2SqlEntity prepareBatchUpdate(Data data) { + H2SqlEntity entity = new H2SqlEntity(); Map source = new HashMap<>(); source.put(InstanceTable.COLUMN_HEARTBEAT_TIME, data.getDataLong(0)); - return source; + String sql = getBatchUpdateSql(InstanceTable.TABLE, source.keySet(), InstanceTable.COLUMN_APPLICATION_ID); + entity.setSql(sql); + List params = new ArrayList<>(source.values()); + params.add(data.getDataString(0)); + entity.setParams(params.toArray(new Object[0])); + return entity; } } diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memory/dao/MemoryMetricH2DAO.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memory/dao/MemoryMetricH2DAO.java index 3afaeaac8..1f15524a3 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memory/dao/MemoryMetricH2DAO.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memory/dao/MemoryMetricH2DAO.java @@ -4,6 +4,7 @@ import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.storage.define.jvm.MemoryMetricTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.storage.h2.define.H2SqlEntity; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import java.util.HashMap; @@ -12,13 +13,15 @@ import java.util.Map; /** * @author pengys5 */ -public class MemoryMetricH2DAO extends H2DAO implements IMemoryMetricDAO, IPersistenceDAO, Map> { +public class MemoryMetricH2DAO extends H2DAO implements IMemoryMetricDAO, IPersistenceDAO { @Override public Data get(String id, DataDefine dataDefine) { return null; } - @Override public Map prepareBatchInsert(Data data) { + @Override public H2SqlEntity prepareBatchInsert(Data data) { + H2SqlEntity entity = new H2SqlEntity(); Map source = new HashMap<>(); + source.put("id", data.getDataString(0)); source.put(MemoryMetricTable.COLUMN_APPLICATION_INSTANCE_ID, data.getDataInteger(0)); source.put(MemoryMetricTable.COLUMN_IS_HEAP, data.getDataBoolean(0)); source.put(MemoryMetricTable.COLUMN_INIT, data.getDataLong(0)); @@ -27,10 +30,13 @@ public class MemoryMetricH2DAO extends H2DAO implements IMemoryMetricDAO, IPersi source.put(MemoryMetricTable.COLUMN_COMMITTED, data.getDataLong(3)); source.put(MemoryMetricTable.COLUMN_TIME_BUCKET, data.getDataLong(4)); - return source; + String sql = getBatchInsertSql(MemoryMetricTable.TABLE, source.keySet()); + entity.setSql(sql); + entity.setParams(source.values().toArray(new Object[0])); + return entity; } - @Override public Map prepareBatchUpdate(Data data) { + @Override public H2SqlEntity prepareBatchUpdate(Data data) { return null; } } diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memorypool/dao/MemoryPoolMetricH2DAO.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memorypool/dao/MemoryPoolMetricH2DAO.java index 7cbd236bb..64a8638ba 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memorypool/dao/MemoryPoolMetricH2DAO.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memorypool/dao/MemoryPoolMetricH2DAO.java @@ -4,6 +4,7 @@ import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.storage.define.jvm.MemoryPoolMetricTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.storage.h2.define.H2SqlEntity; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import java.util.HashMap; @@ -12,13 +13,15 @@ import java.util.Map; /** * @author pengys5 */ -public class MemoryPoolMetricH2DAO extends H2DAO implements IMemoryPoolMetricDAO, IPersistenceDAO, Map> { +public class MemoryPoolMetricH2DAO extends H2DAO implements IMemoryPoolMetricDAO, IPersistenceDAO { @Override public Data get(String id, DataDefine dataDefine) { return null; } - @Override public Map prepareBatchInsert(Data data) { + @Override public H2SqlEntity prepareBatchInsert(Data data) { + H2SqlEntity entity = new H2SqlEntity(); Map source = new HashMap<>(); + source.put("id", data.getDataString(0)); source.put(MemoryPoolMetricTable.COLUMN_INSTANCE_ID, data.getDataInteger(0)); source.put(MemoryPoolMetricTable.COLUMN_POOL_TYPE, data.getDataInteger(1)); source.put(MemoryPoolMetricTable.COLUMN_INIT, data.getDataLong(0)); @@ -27,10 +30,13 @@ public class MemoryPoolMetricH2DAO extends H2DAO implements IMemoryPoolMetricDAO source.put(MemoryPoolMetricTable.COLUMN_COMMITTED, data.getDataLong(3)); source.put(MemoryPoolMetricTable.COLUMN_TIME_BUCKET, data.getDataLong(4)); - return source; + String sql = getBatchInsertSql(MemoryPoolMetricTable.TABLE, source.keySet()); + entity.setSql(sql); + entity.setParams(source.values().toArray(new Object[0])); + return entity; } - @Override public Map prepareBatchUpdate(Data data) { + @Override public H2SqlEntity prepareBatchUpdate(Data data) { return null; } } diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/dao/ApplicationH2DAO.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/dao/ApplicationH2DAO.java index 53337184f..6f45a62cb 100644 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/dao/ApplicationH2DAO.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/dao/ApplicationH2DAO.java @@ -8,12 +8,15 @@ import org.skywalking.apm.collector.storage.h2.dao.H2DAO; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.text.MessageFormat; + /** * @author pengys5 */ public class ApplicationH2DAO extends H2DAO implements IApplicationDAO { private final Logger logger = LoggerFactory.getLogger(ApplicationH2DAO.class); - @Override + private static final String INSERT_APPLICATION_SQL = "insert into {0}({1}, {2}) values(?, ?)"; +; @Override public int getApplicationId(String applicationCode) { logger.info("get the application id with application code = {}", applicationCode); String sql = "select " + ApplicationTable.COLUMN_APPLICATION_ID + " from " + @@ -34,12 +37,12 @@ public class ApplicationH2DAO extends H2DAO implements IApplicationDAO { @Override public void save(ApplicationDataDefine.Application application) { - String insertSQL = "insert into " + ApplicationTable.TABLE + "(" + ApplicationTable.COLUMN_APPLICATION_ID + - "," + ApplicationTable.COLUMN_APPLICATION_CODE + ") values (?, ?)"; H2Client client = getClient(); + String sql = MessageFormat.format(INSERT_APPLICATION_SQL, ApplicationTable.TABLE, ApplicationTable.COLUMN_APPLICATION_ID, + ApplicationTable.COLUMN_APPLICATION_CODE); Object[] params = new Object[] {application.getApplicationId(), application.getApplicationCode()}; try { - client.execute(insertSQL, params); + client.execute(sql, params); } catch (H2ClientException e) { logger.error(e.getMessage(), e); } diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/dao/InstanceH2DAO.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/dao/InstanceH2DAO.java index c11920a7a..6b7ccc612 100644 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/dao/InstanceH2DAO.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/dao/InstanceH2DAO.java @@ -1,33 +1,97 @@ package org.skywalking.apm.collector.agentregister.worker.instance.dao; +import org.skywalking.apm.collector.client.h2.H2Client; +import org.skywalking.apm.collector.client.h2.H2ClientException; +import org.skywalking.apm.collector.core.stream.Data; +import org.skywalking.apm.collector.storage.define.jvm.CpuMetricTable; +import org.skywalking.apm.collector.storage.define.register.ApplicationTable; import org.skywalking.apm.collector.storage.define.register.InstanceDataDefine; +import org.skywalking.apm.collector.storage.define.register.InstanceTable; +import org.skywalking.apm.collector.storage.define.serviceref.ServiceReferenceTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.sql.ResultSet; +import java.sql.SQLException; +import java.text.MessageFormat; +import java.util.HashMap; +import java.util.Map; /** * @author pengys5 */ public class InstanceH2DAO extends H2DAO implements IInstanceDAO { + private final Logger logger = LoggerFactory.getLogger(InstanceH2DAO.class); + private static final String GET_INSTANCE_ID_SQL = "select {0} from {1} where {2} = ? and {3} = ?"; + private static final String GET_APPLICATION_ID_SQL = "select {0} from {1} where {2} = ?"; + private static final String UPDATE_HEARTBEAT_TIME_SQL = "updatte {0} set {1} = ? where {2} = ?"; @Override public int getInstanceId(int applicationId, String agentUUID) { + logger.info("get the application id with application id = {}, agentUUID = {}", applicationId, agentUUID); + H2Client client = getClient(); + String sql = MessageFormat.format(GET_INSTANCE_ID_SQL, InstanceTable.COLUMN_INSTANCE_ID, InstanceTable.TABLE, InstanceTable.COLUMN_APPLICATION_ID, + InstanceTable.COLUMN_AGENT_UUID); + Object[] params = new Object[]{applicationId, agentUUID}; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + return rs.getInt(InstanceTable.COLUMN_INSTANCE_ID); + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } return 0; } @Override public int getMaxInstanceId() { - return 0; + return getMaxId(InstanceTable.TABLE, InstanceTable.COLUMN_INSTANCE_ID); } @Override public int getMinInstanceId() { - return 0; + return getMinId(InstanceTable.TABLE, InstanceTable.COLUMN_INSTANCE_ID); } @Override public void save(InstanceDataDefine.Instance instance) { - + H2Client client = getClient(); + Map source = new HashMap<>(); + source.put(InstanceTable.COLUMN_INSTANCE_ID, instance.getInstanceId()); + 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()); + source.put(InstanceTable.COLUMN_OS_INFO, instance.getOsInfo()); + String sql = getBatchInsertSql(InstanceTable.TABLE, source.keySet()); + Object[] params = source.values().toArray(new Object[0]); + try { + client.execute(sql, params); + } catch (H2ClientException e) { + logger.error(e.getMessage(), e); + } } @Override public void updateHeartbeatTime(int instanceId, long heartbeatTime) { - + H2Client client = getClient(); + String sql = MessageFormat.format(UPDATE_HEARTBEAT_TIME_SQL, InstanceTable.TABLE, InstanceTable.COLUMN_HEARTBEAT_TIME, + InstanceTable.COLUMN_INSTANCE_ID); + Object[] params = new Object[] {heartbeatTime, instanceId}; + try { + client.execute(sql, params); + } catch (H2ClientException e) { + logger.error(e.getMessage(), e); + } } @Override public int getApplicationId(int applicationInstanceId) { + logger.info("get the application id with application id = {}", applicationInstanceId); + H2Client client = getClient(); + String sql = MessageFormat.format(GET_APPLICATION_ID_SQL, InstanceTable.COLUMN_APPLICATION_ID, InstanceTable.TABLE, InstanceTable.COLUMN_APPLICATION_ID); + Object[] params = new Object[]{applicationInstanceId}; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + return rs.getInt(InstanceTable.COLUMN_APPLICATION_ID); + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } return 0; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/dao/GlobalTraceH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/dao/GlobalTraceH2DAO.java index d17002aa0..05dee62ab 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/dao/GlobalTraceH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/dao/GlobalTraceH2DAO.java @@ -1,27 +1,44 @@ package org.skywalking.apm.collector.agentstream.worker.global.dao; -import org.skywalking.apm.collector.agentstream.worker.instance.performance.dao.InstPerformanceH2DAO; +import org.skywalking.apm.collector.core.framework.UnexpectedException; import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.storage.define.DataDefine; +import org.skywalking.apm.collector.storage.define.global.GlobalTraceTable; +import org.skywalking.apm.collector.storage.define.node.NodeComponentTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.storage.h2.define.H2SqlEntity; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.HashMap; import java.util.Map; /** * @author pengys5 */ -public class GlobalTraceH2DAO extends H2DAO implements IGlobalTraceDAO, IPersistenceDAO, Map> { +public class GlobalTraceH2DAO extends H2DAO implements IGlobalTraceDAO, IPersistenceDAO { private final Logger logger = LoggerFactory.getLogger(GlobalTraceH2DAO.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 Map prepareBatchInsert(Data data) { - return null; + + @Override public H2SqlEntity prepareBatchUpdate(Data data) { + throw new UnexpectedException("There is no need to merge stream data with database data."); } - @Override public Map prepareBatchUpdate(Data data) { - return null; + + @Override public H2SqlEntity prepareBatchInsert(Data data) { + Map source = new HashMap<>(); + H2SqlEntity entity = new H2SqlEntity(); + source.put("id", data.getDataString(0)); + source.put(GlobalTraceTable.COLUMN_SEGMENT_ID, data.getDataString(1)); + source.put(GlobalTraceTable.COLUMN_GLOBAL_TRACE_ID, data.getDataString(2)); + source.put(GlobalTraceTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + logger.debug("global trace source: {}", source.toString()); + + String sql = getBatchInsertSql(NodeComponentTable.TABLE, source.keySet()); + entity.setSql(sql); + entity.setParams(source.values().toArray(new Object[0])); + return entity; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/dao/InstPerformanceH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/dao/InstPerformanceH2DAO.java index bb7b9fb04..418191320 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/dao/InstPerformanceH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/dao/InstPerformanceH2DAO.java @@ -1,27 +1,74 @@ package org.skywalking.apm.collector.agentstream.worker.instance.performance.dao; -import org.skywalking.apm.collector.agentstream.worker.node.component.dao.NodeComponentH2DAO; +import org.skywalking.apm.collector.client.h2.H2Client; +import org.skywalking.apm.collector.client.h2.H2ClientException; import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.storage.define.DataDefine; +import org.skywalking.apm.collector.storage.define.instance.InstPerformanceTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.storage.h2.define.H2SqlEntity; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.util.Map; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.text.MessageFormat; +import java.util.*; /** * @author pengys5 */ -public class InstPerformanceH2DAO extends H2DAO implements IInstPerformanceDAO, IPersistenceDAO, Map> { +public class InstPerformanceH2DAO extends H2DAO implements IInstPerformanceDAO, IPersistenceDAO { private final Logger logger = LoggerFactory.getLogger(InstPerformanceH2DAO.class); + private static final String GET_SQL = "select * from {0} where {1} = ?"; @Override public Data get(String id, DataDefine dataDefine) { + H2Client client = getClient(); + String sql = MessageFormat.format(GET_SQL, InstPerformanceTable.TABLE, "id"); + Object[] params = new Object[]{id}; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + Data data = dataDefine.build(id); + data.setDataInteger(0, rs.getInt(InstPerformanceTable.COLUMN_APPLICATION_ID)); + data.setDataInteger(1, rs.getInt(InstPerformanceTable.COLUMN_INSTANCE_ID)); + data.setDataInteger(2, rs.getInt(InstPerformanceTable.COLUMN_CALLS)); + data.setDataLong(0, rs.getLong(InstPerformanceTable.COLUMN_COST_TOTAL)); + data.setDataLong(1, rs.getLong(InstPerformanceTable.COLUMN_TIME_BUCKET)); + return data; + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } return null; } - @Override public Map prepareBatchInsert(Data data) { - return null; + @Override public H2SqlEntity prepareBatchInsert(Data data) { + Map source = new HashMap<>(); + H2SqlEntity entity = new H2SqlEntity(); + source.put("id", data.getDataString(0)); + source.put(InstPerformanceTable.COLUMN_APPLICATION_ID, data.getDataInteger(0)); + source.put(InstPerformanceTable.COLUMN_INSTANCE_ID, data.getDataInteger(1)); + source.put(InstPerformanceTable.COLUMN_CALLS, data.getDataInteger(2)); + source.put(InstPerformanceTable.COLUMN_COST_TOTAL, data.getDataLong(0)); + source.put(InstPerformanceTable.COLUMN_TIME_BUCKET, data.getDataLong(1)); + String sql = getBatchInsertSql(InstPerformanceTable.TABLE, source.keySet()); + entity.setSql(sql); + entity.setParams(source.values().toArray(new Object[0])); + return entity; } - @Override public Map prepareBatchUpdate(Data data) { - return null; + @Override public H2SqlEntity prepareBatchUpdate(Data data) { + Map source = new HashMap<>(); + H2SqlEntity entity = new H2SqlEntity(); + source.put(InstPerformanceTable.COLUMN_APPLICATION_ID, data.getDataInteger(0)); + source.put(InstPerformanceTable.COLUMN_INSTANCE_ID, data.getDataInteger(1)); + source.put(InstPerformanceTable.COLUMN_CALLS, data.getDataInteger(2)); + source.put(InstPerformanceTable.COLUMN_COST_TOTAL, data.getDataLong(0)); + source.put(InstPerformanceTable.COLUMN_TIME_BUCKET, data.getDataLong(1)); + String id = data.getDataString(0); + String sql = getBatchUpdateSql(InstPerformanceTable.TABLE, source.keySet(), "id"); + entity.setSql(sql); + List values = new ArrayList<>(source.values()); + values.add(id); + entity.setParams(values.toArray(new Object[0])); + return entity; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentH2DAO.java index 870f656e6..98f0797f7 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/dao/NodeComponentH2DAO.java @@ -1,26 +1,82 @@ package org.skywalking.apm.collector.agentstream.worker.node.component.dao; +import org.skywalking.apm.collector.client.h2.H2Client; +import org.skywalking.apm.collector.client.h2.H2ClientException; import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.storage.define.DataDefine; +import org.skywalking.apm.collector.storage.define.node.NodeComponentTable; +import org.skywalking.apm.collector.storage.define.serviceref.ServiceReferenceTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.storage.h2.define.H2SqlEntity; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.util.Map; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.text.MessageFormat; +import java.util.*; /** * @author pengys5 */ -public class NodeComponentH2DAO extends H2DAO implements INodeComponentDAO, IPersistenceDAO, Map> { +public class NodeComponentH2DAO extends H2DAO implements INodeComponentDAO, IPersistenceDAO { private final Logger logger = LoggerFactory.getLogger(NodeComponentH2DAO.class); - @Override public Data get(String id, DataDefine dataDefine) { + private static final String GET_SQL = "select * from {0} where {1} = ?"; + + @Override + public Data get(String id, DataDefine dataDefine) { + H2Client client = getClient(); + String sql = MessageFormat.format(GET_SQL, ServiceReferenceTable.TABLE, "id"); + Object[] params = new Object[]{id}; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + Data data = dataDefine.build(id); + data.setDataInteger(0, rs.getInt(NodeComponentTable.COLUMN_COMPONENT_ID)); + data.setDataString(1, rs.getString(NodeComponentTable.COLUMN_COMPONENT_NAME)); + data.setDataInteger(1, rs.getInt(NodeComponentTable.COLUMN_PEER_ID)); + data.setDataString(2, rs.getString(NodeComponentTable.COLUMN_PEER)); + data.setDataLong(0, rs.getLong(NodeComponentTable.COLUMN_TIME_BUCKET)); + return data; + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } return null; } - @Override public Map prepareBatchInsert(Data data) { - return null; + + @Override + public H2SqlEntity prepareBatchInsert(Data data) { + Map source = new HashMap<>(); + H2SqlEntity entity = new H2SqlEntity(); + source.put("id", data.getDataString(0)); + source.put(NodeComponentTable.COLUMN_COMPONENT_ID, data.getDataInteger(0)); + source.put(NodeComponentTable.COLUMN_COMPONENT_NAME, data.getDataString(1)); + source.put(NodeComponentTable.COLUMN_PEER_ID, data.getDataInteger(1)); + source.put(NodeComponentTable.COLUMN_PEER, data.getDataString(2)); + source.put(NodeComponentTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + + String sql = getBatchInsertSql(NodeComponentTable.TABLE, source.keySet()); + entity.setSql(sql); + entity.setParams(source.values().toArray(new Object[0])); + return entity; } - @Override public Map prepareBatchUpdate(Data data) { - return null; + + @Override + public H2SqlEntity prepareBatchUpdate(Data data) { + Map source = new HashMap<>(); + H2SqlEntity entity = new H2SqlEntity(); + source.put(NodeComponentTable.COLUMN_COMPONENT_ID, data.getDataInteger(0)); + source.put(NodeComponentTable.COLUMN_COMPONENT_NAME, data.getDataString(1)); + source.put(NodeComponentTable.COLUMN_PEER_ID, data.getDataInteger(1)); + source.put(NodeComponentTable.COLUMN_PEER, data.getDataString(2)); + source.put(NodeComponentTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + String id = data.getDataString(0); + String sql = getBatchUpdateSql(NodeComponentTable.TABLE, source.keySet(), "id"); + entity.setSql(sql); + List values = new ArrayList<>(source.values()); + values.add(id); + entity.setParams(values.toArray(new Object[0])); + return entity; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingH2DAO.java index e5cb8709a..2f8712c3d 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/dao/NodeMappingH2DAO.java @@ -1,27 +1,72 @@ package org.skywalking.apm.collector.agentstream.worker.node.mapping.dao; import org.skywalking.apm.collector.agentstream.worker.node.component.dao.NodeComponentH2DAO; +import org.skywalking.apm.collector.client.h2.H2Client; +import org.skywalking.apm.collector.client.h2.H2ClientException; import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.storage.define.DataDefine; +import org.skywalking.apm.collector.storage.define.node.NodeMappingTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.storage.h2.define.H2SqlEntity; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.util.Map; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.text.MessageFormat; +import java.util.*; /** * @author pengys5 */ -public class NodeMappingH2DAO extends H2DAO implements INodeMappingDAO, IPersistenceDAO, Map> { +public class NodeMappingH2DAO extends H2DAO implements INodeMappingDAO, IPersistenceDAO { private final Logger logger = LoggerFactory.getLogger(NodeComponentH2DAO.class); + private static final String GET_SQL = "select * from {0} where {1} = ?"; @Override public Data get(String id, DataDefine dataDefine) { + H2Client client = getClient(); + String sql = MessageFormat.format(GET_SQL, NodeMappingTable.TABLE, "id"); + Object[] params = new Object[]{id}; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + Data data = dataDefine.build(id); + data.setDataInteger(0, rs.getInt(NodeMappingTable.COLUMN_APPLICATION_ID)); + data.setDataInteger(1, rs.getInt(NodeMappingTable.COLUMN_ADDRESS_ID)); + data.setDataString(1, rs.getString(NodeMappingTable.COLUMN_ADDRESS)); + data.setDataLong(0, rs.getLong(NodeMappingTable.COLUMN_TIME_BUCKET)); + return data; + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } return null; } - @Override public Map prepareBatchInsert(Data data) { - return null; + @Override public H2SqlEntity prepareBatchInsert(Data data) { + Map source = new HashMap<>(); + H2SqlEntity entity = new H2SqlEntity(); + source.put("id", data.getDataString(0)); + source.put(NodeMappingTable.COLUMN_APPLICATION_ID, data.getDataInteger(0)); + source.put(NodeMappingTable.COLUMN_ADDRESS_ID, data.getDataInteger(1)); + source.put(NodeMappingTable.COLUMN_ADDRESS, data.getDataString(1)); + source.put(NodeMappingTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + String sql = getBatchInsertSql(NodeMappingTable.TABLE, source.keySet()); + entity.setSql(sql); + entity.setParams(source.values().toArray(new Object[0])); + return entity; } - @Override public Map prepareBatchUpdate(Data data) { - return null; + @Override public H2SqlEntity prepareBatchUpdate(Data data) { + Map source = new HashMap<>(); + H2SqlEntity entity = new H2SqlEntity(); + source.put(NodeMappingTable.COLUMN_APPLICATION_ID, data.getDataInteger(0)); + source.put(NodeMappingTable.COLUMN_ADDRESS_ID, data.getDataInteger(1)); + source.put(NodeMappingTable.COLUMN_ADDRESS, data.getDataString(1)); + source.put(NodeMappingTable.COLUMN_TIME_BUCKET, data.getDataLong(0)); + String id = data.getDataString(0); + String sql = getBatchUpdateSql(NodeMappingTable.TABLE, source.keySet(), "id"); + entity.setSql(sql); + List values = new ArrayList<>(source.values()); + values.add(id); + entity.setParams(values.toArray(new Object[0])); + return entity; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/dao/NodeReferenceH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/dao/NodeReferenceH2DAO.java index 7fe20ff63..f03932c9f 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/dao/NodeReferenceH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/dao/NodeReferenceH2DAO.java @@ -1,27 +1,25 @@ package org.skywalking.apm.collector.agentstream.worker.noderef.dao; -import org.skywalking.apm.collector.agentstream.worker.segment.cost.dao.SegmentCostH2DAO; import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.storage.h2.define.H2SqlEntity; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.util.Map; - /** * @author pengys5 */ -public class NodeReferenceH2DAO extends H2DAO implements INodeReferenceDAO, IPersistenceDAO, Map> { +public class NodeReferenceH2DAO extends H2DAO implements INodeReferenceDAO, IPersistenceDAO { private final Logger logger = LoggerFactory.getLogger(NodeReferenceH2DAO.class); @Override public Data get(String id, DataDefine dataDefine) { return null; } - @Override public Map prepareBatchInsert(Data data) { + @Override public H2SqlEntity prepareBatchInsert(Data data) { return null; } - @Override public Map prepareBatchUpdate(Data data) { + @Override public H2SqlEntity prepareBatchUpdate(Data data) { return null; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/SegmentCostH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/SegmentCostH2DAO.java index 1fb2e7181..77f9f2628 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/SegmentCostH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/dao/SegmentCostH2DAO.java @@ -1,34 +1,46 @@ package org.skywalking.apm.collector.agentstream.worker.segment.cost.dao; -import org.skywalking.apm.collector.agentstream.worker.service.entry.dao.ServiceEntryH2DAO; -import org.skywalking.apm.collector.client.h2.H2Client; -import org.skywalking.apm.collector.client.h2.H2ClientException; import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.storage.define.DataDefine; -import org.skywalking.apm.collector.storage.define.service.ServiceEntryTable; -import org.skywalking.apm.collector.storage.define.serviceref.ServiceReferenceTable; +import org.skywalking.apm.collector.storage.define.segment.SegmentCostTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.storage.h2.define.H2SqlEntity; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.sql.ResultSet; -import java.sql.SQLException; import java.util.HashMap; import java.util.Map; /** * @author pengys5 */ -public class SegmentCostH2DAO extends H2DAO implements ISegmentCostDAO, IPersistenceDAO, Map> { +public class SegmentCostH2DAO extends H2DAO implements ISegmentCostDAO, IPersistenceDAO { private final Logger logger = LoggerFactory.getLogger(SegmentCostH2DAO.class); @Override public Data get(String id, DataDefine dataDefine) { return null; } - @Override public Map prepareBatchInsert(Data data) { - return null; + @Override public H2SqlEntity prepareBatchInsert(Data data) { + logger.debug("segment cost prepareBatchInsert, id: {}", data.getDataString(0)); + H2SqlEntity entity = new H2SqlEntity(); + Map source = new HashMap<>(); + source.put("id", data.getDataString(0)); + source.put(SegmentCostTable.COLUMN_SEGMENT_ID, data.getDataString(1)); + source.put(SegmentCostTable.COLUMN_APPLICATION_ID, data.getDataInteger(0)); + source.put(SegmentCostTable.COLUMN_SERVICE_NAME, data.getDataString(2)); + source.put(SegmentCostTable.COLUMN_COST, data.getDataLong(0)); + source.put(SegmentCostTable.COLUMN_START_TIME, data.getDataLong(1)); + source.put(SegmentCostTable.COLUMN_END_TIME, data.getDataLong(2)); + source.put(SegmentCostTable.COLUMN_IS_ERROR, data.getDataBoolean(0)); + source.put(SegmentCostTable.COLUMN_TIME_BUCKET, data.getDataLong(3)); + logger.debug("segment cost source: {}", source.toString()); + + String sql = getBatchInsertSql(SegmentCostTable.TABLE, source.keySet()); + entity.setSql(sql); + entity.setParams(source.values().toArray(new Object[0])); + return entity; } - @Override public Map prepareBatchUpdate(Data data) { + @Override public H2SqlEntity prepareBatchUpdate(Data data) { return null; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/SegmentH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/SegmentH2DAO.java index 37054c67b..1d78a8434 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/SegmentH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/dao/SegmentH2DAO.java @@ -3,25 +3,38 @@ package org.skywalking.apm.collector.agentstream.worker.segment.origin.dao; import org.skywalking.apm.collector.agentstream.worker.segment.cost.dao.SegmentCostH2DAO; import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.storage.define.DataDefine; +import org.skywalking.apm.collector.storage.define.segment.SegmentTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.storage.h2.define.H2SqlEntity; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.Base64; +import java.util.HashMap; import java.util.Map; /** * @author pengys5 */ -public class SegmentH2DAO extends H2DAO implements ISegmentDAO, IPersistenceDAO, Map> { +public class SegmentH2DAO extends H2DAO implements ISegmentDAO, IPersistenceDAO { private final Logger logger = LoggerFactory.getLogger(SegmentCostH2DAO.class); @Override public Data get(String id, DataDefine dataDefine) { return null; } - @Override public Map prepareBatchInsert(Data data) { - return null; + @Override public H2SqlEntity prepareBatchInsert(Data data) { + Map source = new HashMap<>(); + H2SqlEntity entity = new H2SqlEntity(); + source.put("id", data.getDataString(0)); + source.put(SegmentTable.COLUMN_DATA_BINARY, new String(Base64.getEncoder().encode(data.getDataBytes(0)))); + logger.debug("segment source: {}", source.toString()); + + String sql = getBatchInsertSql(SegmentTable.TABLE, source.keySet()); + entity.setSql(sql); + entity.setParams(source.values().toArray(new Object[0])); + return entity; } - @Override public Map prepareBatchUpdate(Data data) { + @Override public H2SqlEntity prepareBatchUpdate(Data data) { return null; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/ServiceEntryH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/ServiceEntryH2DAO.java index a2650151b..aa182958e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/ServiceEntryH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/dao/ServiceEntryH2DAO.java @@ -1,6 +1,5 @@ package org.skywalking.apm.collector.agentstream.worker.service.entry.dao; -import org.skywalking.apm.collector.agentstream.worker.serviceref.dao.ServiceReferenceH2DAO; import org.skywalking.apm.collector.client.h2.H2Client; import org.skywalking.apm.collector.client.h2.H2ClientException; import org.skywalking.apm.collector.core.stream.Data; @@ -8,27 +7,25 @@ import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.storage.define.service.ServiceEntryTable; import org.skywalking.apm.collector.storage.define.serviceref.ServiceReferenceTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.storage.h2.define.H2SqlEntity; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.sql.ResultSet; import java.sql.SQLException; -import java.util.HashMap; -import java.util.Map; +import java.util.*; /** * @author pengys5 */ -public class ServiceEntryH2DAO extends H2DAO implements IServiceEntryDAO, IPersistenceDAO, Map> { +public class ServiceEntryH2DAO extends H2DAO implements IServiceEntryDAO, IPersistenceDAO { private final Logger logger = LoggerFactory.getLogger(ServiceEntryH2DAO.class); @Override public Data get(String id, DataDefine dataDefine) { H2Client client = getClient(); - String sql = "select * from " + ServiceReferenceTable.TABLE + " where " + ServiceReferenceTable.COLUMN_ENTRY_SERVICE_ID + "=?"; + String sql = "select * from " + ServiceReferenceTable.TABLE + " where id = ?"; Object[] params = new Object[] {id}; - ResultSet rs = null; - try { - rs = client.executeQuery(sql, params); + try (ResultSet rs = client.executeQuery(sql, params)) { Data data = dataDefine.build(id); data.setDataInteger(0, rs.getInt(ServiceEntryTable.COLUMN_APPLICATION_ID)); data.setDataInteger(1, rs.getInt(ServiceEntryTable.COLUMN_ENTRY_SERVICE_ID)); @@ -38,27 +35,37 @@ public class ServiceEntryH2DAO extends H2DAO implements IServiceEntryDAO, IPersi return data; } catch (SQLException | H2ClientException e) { logger.error(e.getMessage(), e); - } finally { - client.closeResultSet(rs); } return null; } - @Override public Map prepareBatchInsert(Data data) { + @Override public H2SqlEntity prepareBatchInsert(Data data) { + H2SqlEntity entity = new H2SqlEntity(); Map source = new HashMap<>(); + source.put("id", data.getDataString(0)); source.put(ServiceEntryTable.COLUMN_APPLICATION_ID, data.getDataInteger(0)); source.put(ServiceEntryTable.COLUMN_ENTRY_SERVICE_ID, data.getDataInteger(1)); source.put(ServiceEntryTable.COLUMN_ENTRY_SERVICE_NAME, data.getDataString(1)); source.put(ServiceEntryTable.COLUMN_REGISTER_TIME, data.getDataLong(0)); source.put(ServiceEntryTable.COLUMN_NEWEST_TIME, data.getDataLong(1)); - return source; + String sql = getBatchInsertSql(ServiceEntryTable.TABLE, source.keySet()); + entity.setSql(sql); + entity.setParams(source.values().toArray(new Object[0])); + return entity; } - @Override public Map prepareBatchUpdate(Data data) { + @Override public H2SqlEntity prepareBatchUpdate(Data data) { + H2SqlEntity entity = new H2SqlEntity(); Map source = new HashMap<>(); source.put(ServiceEntryTable.COLUMN_APPLICATION_ID, data.getDataInteger(0)); source.put(ServiceEntryTable.COLUMN_ENTRY_SERVICE_ID, data.getDataInteger(1)); source.put(ServiceEntryTable.COLUMN_ENTRY_SERVICE_NAME, data.getDataString(1)); source.put(ServiceEntryTable.COLUMN_REGISTER_TIME, data.getDataLong(0)); source.put(ServiceEntryTable.COLUMN_NEWEST_TIME, data.getDataLong(1)); - return source; + String id = data.getDataString(0); + String sql = getBatchUpdateSql(ServiceEntryTable.TABLE, source.keySet(), "id"); + entity.setSql(sql); + List values = new ArrayList<>(source.values()); + values.add(id); + entity.setParams(values.toArray(new Object[0])); + return entity; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/dao/ServiceReferenceH2DAO.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/dao/ServiceReferenceH2DAO.java index 6b17bf8ec..5cf648dea 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/dao/ServiceReferenceH2DAO.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/dao/ServiceReferenceH2DAO.java @@ -2,9 +2,9 @@ package org.skywalking.apm.collector.agentstream.worker.serviceref.dao; import org.skywalking.apm.collector.client.h2.H2Client; import org.skywalking.apm.collector.client.h2.H2ClientException; -import org.skywalking.apm.collector.storage.define.register.InstanceTable; import org.skywalking.apm.collector.storage.define.serviceref.ServiceReferenceTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.storage.h2.define.H2SqlEntity; import org.skywalking.apm.collector.stream.worker.impl.dao.IPersistenceDAO; import org.skywalking.apm.collector.core.stream.Data; import org.skywalking.apm.collector.storage.define.DataDefine; @@ -13,47 +13,50 @@ import org.slf4j.LoggerFactory; import java.sql.ResultSet; import java.sql.SQLException; -import java.util.HashMap; -import java.util.Map; +import java.text.MessageFormat; +import java.util.*; /** * @author pengys5 */ -public class ServiceReferenceH2DAO extends H2DAO implements IServiceReferenceDAO, IPersistenceDAO, Map> { +public class ServiceReferenceH2DAO extends H2DAO implements IServiceReferenceDAO, IPersistenceDAO { private final Logger logger = LoggerFactory.getLogger(ServiceReferenceH2DAO.class); - @Override public Data get(String id, DataDefine dataDefine) { + private static final String GET_SQL = "select * from {0} where {1} = ?"; + @Override + public Data get(String id, DataDefine dataDefine) { H2Client client = getClient(); - String sql = "select * from " + ServiceReferenceTable.TABLE + " where " + ServiceReferenceTable.COLUMN_ENTRY_SERVICE_ID + "=?"; - Object[] params = new Object[] {id}; - ResultSet rs = null; - try { - rs = client.executeQuery(sql, params); - Data data = dataDefine.build(id); - data.setDataInteger(0, rs.getInt(ServiceReferenceTable.COLUMN_ENTRY_SERVICE_ID)); - data.setDataString(1, rs.getString(ServiceReferenceTable.COLUMN_ENTRY_SERVICE_NAME)); - data.setDataInteger(1, rs.getInt(ServiceReferenceTable.COLUMN_FRONT_SERVICE_ID)); - data.setDataString(2, rs.getString(ServiceReferenceTable.COLUMN_FRONT_SERVICE_NAME)); - data.setDataInteger(2, rs.getInt(ServiceReferenceTable.COLUMN_BEHIND_SERVICE_ID)); - data.setDataString(3, rs.getString(ServiceReferenceTable.COLUMN_BEHIND_SERVICE_NAME)); - data.setDataLong(0, rs.getLong(ServiceReferenceTable.COLUMN_S1_LTE)); - data.setDataLong(1, rs.getLong(ServiceReferenceTable.COLUMN_S3_LTE)); - data.setDataLong(2, rs.getLong(ServiceReferenceTable.COLUMN_S5_LTE)); - data.setDataLong(3, rs.getLong(ServiceReferenceTable.COLUMN_S5_GT)); - data.setDataLong(4, rs.getLong(ServiceReferenceTable.COLUMN_SUMMARY)); - data.setDataLong(5, rs.getLong(ServiceReferenceTable.COLUMN_ERROR)); - data.setDataLong(6, rs.getLong(ServiceReferenceTable.COLUMN_COST_SUMMARY)); - data.setDataLong(7, rs.getLong(ServiceReferenceTable.COLUMN_TIME_BUCKET)); - return data; + String sql = MessageFormat.format(GET_SQL, ServiceReferenceTable.TABLE, "id"); + Object[] params = new Object[]{id}; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + Data data = dataDefine.build(id); + data.setDataInteger(0, rs.getInt(ServiceReferenceTable.COLUMN_ENTRY_SERVICE_ID)); + data.setDataString(1, rs.getString(ServiceReferenceTable.COLUMN_ENTRY_SERVICE_NAME)); + data.setDataInteger(1, rs.getInt(ServiceReferenceTable.COLUMN_FRONT_SERVICE_ID)); + data.setDataString(2, rs.getString(ServiceReferenceTable.COLUMN_FRONT_SERVICE_NAME)); + data.setDataInteger(2, rs.getInt(ServiceReferenceTable.COLUMN_BEHIND_SERVICE_ID)); + data.setDataString(3, rs.getString(ServiceReferenceTable.COLUMN_BEHIND_SERVICE_NAME)); + data.setDataLong(0, rs.getLong(ServiceReferenceTable.COLUMN_S1_LTE)); + data.setDataLong(1, rs.getLong(ServiceReferenceTable.COLUMN_S3_LTE)); + data.setDataLong(2, rs.getLong(ServiceReferenceTable.COLUMN_S5_LTE)); + data.setDataLong(3, rs.getLong(ServiceReferenceTable.COLUMN_S5_GT)); + data.setDataLong(4, rs.getLong(ServiceReferenceTable.COLUMN_SUMMARY)); + data.setDataLong(5, rs.getLong(ServiceReferenceTable.COLUMN_ERROR)); + data.setDataLong(6, rs.getLong(ServiceReferenceTable.COLUMN_COST_SUMMARY)); + data.setDataLong(7, rs.getLong(ServiceReferenceTable.COLUMN_TIME_BUCKET)); + return data; + } } catch (SQLException | H2ClientException e) { logger.error(e.getMessage(), e); - } finally { - client.closeResultSet(rs); } return null; } - @Override public Map prepareBatchInsert(Data data) { + @Override + public H2SqlEntity prepareBatchInsert(Data data) { + H2SqlEntity entity = new H2SqlEntity(); Map source = new HashMap<>(); + source.put("id", data.getDataString(0)); source.put(ServiceReferenceTable.COLUMN_ENTRY_SERVICE_ID, data.getDataInteger(0)); source.put(ServiceReferenceTable.COLUMN_ENTRY_SERVICE_NAME, data.getDataString(1)); source.put(ServiceReferenceTable.COLUMN_FRONT_SERVICE_ID, data.getDataInteger(1)); @@ -69,10 +72,15 @@ public class ServiceReferenceH2DAO extends H2DAO implements IServiceReferenceDAO source.put(ServiceReferenceTable.COLUMN_COST_SUMMARY, data.getDataLong(6)); source.put(ServiceReferenceTable.COLUMN_TIME_BUCKET, data.getDataLong(7)); - return source; + String sql = getBatchInsertSql(ServiceReferenceTable.TABLE, source.keySet()); + entity.setSql(sql); + entity.setParams(source.values().toArray(new Object[0])); + return entity; } - @Override public Map prepareBatchUpdate(Data data) { + @Override + public H2SqlEntity prepareBatchUpdate(Data data) { + H2SqlEntity entity = new H2SqlEntity(); Map source = new HashMap<>(); source.put(ServiceReferenceTable.COLUMN_ENTRY_SERVICE_ID, data.getDataInteger(0)); source.put(ServiceReferenceTable.COLUMN_ENTRY_SERVICE_NAME, data.getDataString(1)); @@ -89,6 +97,12 @@ public class ServiceReferenceH2DAO extends H2DAO implements IServiceReferenceDAO source.put(ServiceReferenceTable.COLUMN_COST_SUMMARY, data.getDataLong(6)); source.put(ServiceReferenceTable.COLUMN_TIME_BUCKET, data.getDataLong(7)); - return source; + String id = data.getDataString(0); + String sql = getBatchUpdateSql(ServiceReferenceTable.TABLE, source.keySet(), "id"); + entity.setSql(sql); + List values = new ArrayList<>(source.values()); + values.add(id); + entity.setParams(values.toArray(new Object[0])); + return entity; } } diff --git a/apm-collector/apm-collector-boot/src/main/resources/application.yml b/apm-collector/apm-collector-boot/src/main/resources/application.yml index e8a4837c7..dd9eefa56 100644 --- a/apm-collector/apm-collector-boot/src/main/resources/application.yml +++ b/apm-collector/apm-collector-boot/src/main/resources/application.yml @@ -4,20 +4,20 @@ cluster: sessionTimeout: 100000 agent_server: jetty: - host: 0.0.0.0 + host: 127.0.0.1 port: 10800 context_path: / agent_stream: grpc: - host: 192.168.10.13 + host: 127.0.0.1 port: 11800 jetty: - host: 0.0.0.0 + host: 127.0.0.1 port: 12800 context_path: / ui: jetty: - host: 0.0.0.0 + host: 127.0.0.1 port: 12800 context_path: / collector_inside: diff --git a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/h2/H2Client.java b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/h2/H2Client.java index 035383716..fb62a1b1b 100644 --- a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/h2/H2Client.java +++ b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/h2/H2Client.java @@ -55,9 +55,7 @@ public class H2Client implements Client { } public void execute(String sql) throws H2ClientException { - Statement statement = null; - try { - statement = getConnection().createStatement(); + try (Statement statement = getConnection().createStatement()) { statement.execute(sql); statement.closeOnCompletion(); } catch (SQLException e) { @@ -67,13 +65,13 @@ public class H2Client implements Client { public ResultSet executeQuery(String sql, Object[] params) throws H2ClientException { logger.info("execute query with result: {}", sql); - PreparedStatement statement; ResultSet rs; + PreparedStatement statement; try { statement = getConnection().prepareStatement(sql); if (params != null) { for (int i = 0; i < params.length; i++) { - statement.setObject(i+1, params[i]); + statement.setObject(i + 1, params[i]); } } rs = statement.executeQuery(); @@ -86,30 +84,20 @@ public class H2Client implements Client { public boolean execute(String sql, Object[] params) throws H2ClientException { logger.info("execute insert/update/delete: {}", sql); - PreparedStatement statement; boolean flag; - try { - statement = getConnection().prepareStatement(sql); + Connection conn = getConnection(); + try (PreparedStatement statement = conn.prepareStatement(sql)) { + conn.setAutoCommit(false); if (params != null) { for (int i = 0; i < params.length; i++) { - statement.setObject(i+1, params[i]); + statement.setObject(i + 1, params[i]); } } flag = statement.execute(); - statement.closeOnCompletion(); + conn.commit(); } catch (SQLException e) { throw new H2ClientException(e.getMessage(), e); } return flag; } - - public void closeResultSet(ResultSet rs) { - if (rs != null) { - try { - rs.close(); - } catch (SQLException e) { - logger.error(e.getMessage(), e); - } - } - } } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/BatchH2DAO.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/BatchH2DAO.java index 27b69f6d4..397f13bbb 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/BatchH2DAO.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/BatchH2DAO.java @@ -1,14 +1,73 @@ package org.skywalking.apm.collector.storage.h2.dao; -import java.util.List; +import org.skywalking.apm.collector.client.h2.H2ClientException; import org.skywalking.apm.collector.storage.dao.IBatchDAO; +import org.skywalking.apm.collector.storage.h2.define.H2SqlEntity; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.SQLException; +import java.util.HashMap; +import java.util.List; +import java.util.Map; /** * @author pengys5 */ public class BatchH2DAO extends H2DAO implements IBatchDAO { + private final Logger logger = LoggerFactory.getLogger(BatchH2DAO.class); + private final Map batchSqls = new HashMap<>(); - @Override public void batchPersistence(List batchCollection) { + @Override + public void batchPersistence(List batchCollection) { + if (batchCollection != null && batchCollection.size() > 0) { + logger.info("the batch collection size is {}", batchCollection.size()); + Connection conn = null; + try { + conn = getClient().getConnection(); + conn.setAutoCommit(false); + PreparedStatement ps; + for (Object entity : batchCollection) { + H2SqlEntity e = getH2SqlEntity(entity); + String sql = e.getSql(); + if (batchSqls.containsKey(sql)) { + ps = batchSqls.get(sql); + } else { + ps = conn.prepareStatement(sql); + batchSqls.put(sql, ps); + } + Object[] params = e.getParams(); + if (params != null) { + logger.info("the sql is {}, params size is {}", e.getSql(), params.length); + for (int i = 0; i < params.length; i++) { + ps.setObject(i + 1, params[i]); + } + } + ps.addBatch(); + } + + for (String k : batchSqls.keySet()) { + batchSqls.get(k).executeBatch(); + } + conn.commit(); + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + try { + conn.rollback(); + } catch (SQLException e1) { + logger.error(e.getMessage(), e1); + } + } + } + } + + private H2SqlEntity getH2SqlEntity(Object entity) { + if (entity instanceof H2SqlEntity) { + return (H2SqlEntity) entity; + } + return null; } } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/H2DAO.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/H2DAO.java index 1d703e954..f094c4063 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/H2DAO.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/H2DAO.java @@ -8,6 +8,7 @@ import org.slf4j.LoggerFactory; import java.sql.ResultSet; import java.sql.SQLException; +import java.util.Set; /** * @author pengys5 @@ -26,9 +27,7 @@ public abstract class H2DAO extends DAO { public final int getIntValueBySQL(String sql) { H2Client client = getClient(); - ResultSet rs = null; - try { - rs = client.executeQuery(sql, null); + try (ResultSet rs = client.executeQuery(sql, null)) { if (rs.next()) { int id = rs.getInt(1); if (id == Integer.MAX_VALUE || id == Integer.MIN_VALUE) { @@ -39,9 +38,34 @@ public abstract class H2DAO extends DAO { } } catch (SQLException | H2ClientException e) { logger.error(e.getMessage(), e); - } finally { - client.closeResultSet(rs); } return 0; } + + public final String getBatchInsertSql(String tableName, Set columnNames) { + StringBuilder sb = new StringBuilder("insert into "); + sb.append(tableName).append("("); + columnNames.forEach((columnName) -> { + sb.append(columnName).append(","); + }); + sb.delete(sb.length() - 1, sb.length()); + sb.append(") values("); + for (int i = 0; i < columnNames.size(); i++) { + sb.append("?,"); + } + sb.delete(sb.length() - 1, sb.length()); + sb.append(")"); + return sb.toString(); + } + + public final String getBatchUpdateSql(String tableName, Set columnNames, String whereClauseName) { + StringBuilder sb = new StringBuilder("update "); + sb.append(tableName).append(" "); + columnNames.forEach((columnName) -> { + sb.append("set ").append(columnName).append("=?,"); + }); + sb.delete(sb.length() - 1, sb.length()); + sb.append(" where ").append(whereClauseName).append("=?"); + return sb.toString(); + } } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2SqlEntity.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2SqlEntity.java new file mode 100644 index 000000000..c36c21d46 --- /dev/null +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2SqlEntity.java @@ -0,0 +1,22 @@ +package org.skywalking.apm.collector.storage.h2.define; + +public class H2SqlEntity { + private String sql; + private Object[] params; + + public String getSql() { + return sql; + } + + public void setSql(String sql) { + this.sql = sql; + } + + public Object[] getParams() { + return params; + } + + public void setParams(Object[] params) { + this.params = params; + } +} diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2StorageInstaller.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2StorageInstaller.java index 12718a66c..20ecf15c8 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2StorageInstaller.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2StorageInstaller.java @@ -1,10 +1,5 @@ package org.skywalking.apm.collector.storage.h2.define; -import java.sql.ResultSet; -import java.sql.SQLException; -import java.util.List; - -import org.h2.util.IOUtils; import org.skywalking.apm.collector.client.h2.H2Client; import org.skywalking.apm.collector.client.h2.H2ClientException; import org.skywalking.apm.collector.core.client.Client; @@ -15,6 +10,10 @@ import org.skywalking.apm.collector.core.storage.TableDefine; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.util.List; + /** * @author pengys5 */ diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ApplicationH2DAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ApplicationH2DAO.java index 7809287a4..149d73202 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ApplicationH2DAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ApplicationH2DAO.java @@ -1,13 +1,36 @@ package org.skywalking.apm.collector.ui.dao; +import org.skywalking.apm.collector.client.h2.H2Client; +import org.skywalking.apm.collector.client.h2.H2ClientException; +import org.skywalking.apm.collector.core.util.Const; +import org.skywalking.apm.collector.storage.define.register.ApplicationTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.sql.ResultSet; +import java.sql.SQLException; +import java.text.MessageFormat; /** * @author pengys5 */ public class ApplicationH2DAO extends H2DAO implements IApplicationDAO { + private final Logger logger = LoggerFactory.getLogger(ApplicationH2DAO.class); + private static final String GET_APPLICATION_CODE_SQL = "select {0} from {1} where {2} = ?"; @Override public String getApplicationCode(int applicationId) { - return null; + logger.debug("get application code, applicationId: {}", applicationId); + H2Client client = getClient(); + String sql = MessageFormat.format(GET_APPLICATION_CODE_SQL, ApplicationTable.COLUMN_APPLICATION_CODE, ApplicationTable.TABLE, ApplicationTable.COLUMN_APPLICATION_ID); + Object[] params = new Object[]{applicationId}; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + return rs.getString(1); + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } + return Const.UNKNOWN; } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/CpuMetricH2DAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/CpuMetricH2DAO.java new file mode 100644 index 000000000..6cd009208 --- /dev/null +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/CpuMetricH2DAO.java @@ -0,0 +1,81 @@ +package org.skywalking.apm.collector.ui.dao; + +import com.google.gson.JsonArray; +import org.elasticsearch.action.get.GetResponse; +import org.elasticsearch.action.get.MultiGetItemResponse; +import org.elasticsearch.action.get.MultiGetRequestBuilder; +import org.elasticsearch.action.get.MultiGetResponse; +import org.skywalking.apm.collector.client.h2.H2Client; +import org.skywalking.apm.collector.client.h2.H2ClientException; +import org.skywalking.apm.collector.core.util.Const; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; +import org.skywalking.apm.collector.storage.define.jvm.CpuMetricTable; +import org.skywalking.apm.collector.storage.define.register.InstanceDataDefine; +import org.skywalking.apm.collector.storage.define.register.InstanceTable; +import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; +import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.sql.ResultSet; +import java.sql.SQLException; +import java.text.MessageFormat; +import java.util.ArrayList; +import java.util.List; + +/** + * @author pengys5 + */ +public class CpuMetricH2DAO extends H2DAO implements ICpuMetricDAO { + private final Logger logger = LoggerFactory.getLogger(InstanceH2DAO.class); + private static final String GET_METRIC_SQL = "select * from {0} where {1} = ?"; + private static final String GET_METRICS_SQL = "select * from {0} where {1} in ("; + @Override public int getMetric(int instanceId, long timeBucket) { + String id = timeBucket + Const.ID_SPLIT + instanceId; + H2Client client = getClient(); + String sql = MessageFormat.format(GET_METRIC_SQL, CpuMetricTable.TABLE, "id"); + Object[] params = new Object[]{id}; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + return rs.getInt(CpuMetricTable.COLUMN_USAGE_PERCENT); + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } + return 0; + } + + @Override public JsonArray getMetric(int instanceId, long startTimeBucket, long endTimeBucket) { + H2Client client = getClient(); + String sql = MessageFormat.format(GET_METRICS_SQL, CpuMetricTable.TABLE, "id"); + + long timeBucket = startTimeBucket; + List idList = new ArrayList<>(); + do { + timeBucket = TimeBucketUtils.INSTANCE.addSecondForSecondTimeBucket(TimeBucketUtils.TimeBucketType.SECOND.name(), timeBucket, 1); + String id = timeBucket + Const.ID_SPLIT + instanceId; + idList.add(id); + } + while (timeBucket <= endTimeBucket); + + StringBuilder builder = new StringBuilder(); + for( int i = 0 ; i < idList.size(); i++ ) { + builder.append("?,"); + } + builder.delete(builder.length() - 1, builder.length()); + builder.append(")"); + sql = sql + builder; + Object[] params = idList.toArray(new String[0]); + + JsonArray metrics = new JsonArray(); + try (ResultSet rs = client.executeQuery(sql, params)) { + while (rs.next()) { + double cpuUsed = rs.getDouble(CpuMetricTable.COLUMN_USAGE_PERCENT); + metrics.add((int)(cpuUsed * 100)); + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } + return metrics; + } +} diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/GCMetricH2DAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/GCMetricH2DAO.java index b72d5d3d4..bf38132df 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/GCMetricH2DAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/GCMetricH2DAO.java @@ -2,14 +2,19 @@ package org.skywalking.apm.collector.ui.dao; import com.google.gson.JsonObject; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author pengys5 */ public class GCMetricH2DAO extends H2DAO implements IGCMetricDAO { - + private final Logger logger = LoggerFactory.getLogger(GCMetricH2DAO.class); + private static final String GET_GC_COUNT_SQL = "select sum({0}) as cnt, {1} from {2} where {3} > ? group by {1}"; @Override public GCCount getGCCount(long[] timeBuckets, int instanceId) { - return null; + GCCount gcCount = new GCCount(); + + return gcCount; } @Override public JsonObject getMetric(int instanceId, long timeBucket) { diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceH2DAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceH2DAO.java index 7dc331d19..c11ac37c9 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceH2DAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/InstPerformanceH2DAO.java @@ -1,30 +1,156 @@ package org.skywalking.apm.collector.ui.dao; import com.google.gson.JsonArray; +import org.skywalking.apm.collector.client.h2.H2Client; +import org.skywalking.apm.collector.client.h2.H2ClientException; +import org.skywalking.apm.collector.core.util.Const; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; +import org.skywalking.apm.collector.storage.define.instance.InstPerformanceTable; +import org.skywalking.apm.collector.storage.define.jvm.CpuMetricTable; +import org.skywalking.apm.collector.storage.define.register.InstanceDataDefine; +import org.skywalking.apm.collector.storage.define.register.InstanceTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.sql.ResultSet; +import java.sql.SQLException; +import java.text.MessageFormat; +import java.util.ArrayList; +import java.util.List; /** * @author pengys5 */ public class InstPerformanceH2DAO extends H2DAO implements IInstPerformanceDAO { - + private final Logger logger = LoggerFactory.getLogger(InstPerformanceH2DAO.class); + private static final String GET_INST_PERF_SQL = "select * from {0} where {1} = ? and {2} in ("; + private static final String GET_TPS_METRIC_SQL = "select * from {0} where {1} = ?"; + private static final String GET_TPS_METRICS_SQL = "select * from {0} where {1} in ("; @Override public InstPerformance get(long[] timeBuckets, int instanceId) { + H2Client client = getClient(); + logger.info("the inst performance inst id = {}", instanceId); + String sql = MessageFormat.format(GET_INST_PERF_SQL, InstPerformanceTable.TABLE, InstPerformanceTable.COLUMN_INSTANCE_ID, InstPerformanceTable.COLUMN_TIME_BUCKET); + StringBuilder builder = new StringBuilder(); + for( int i = 0 ; i < timeBuckets.length; i++ ) { + builder.append("?,"); + } + builder.delete(builder.length() - 1, builder.length()); + builder.append(")"); + sql = sql + builder; + Object[] params = new Object[timeBuckets.length + 1]; + for(int i = 0; i < timeBuckets.length; i++) { + params[i + 1] = timeBuckets[i]; + } + params[0] = instanceId; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + int callTimes = rs.getInt(InstPerformanceTable.COLUMN_CALLS); + int costTotal = rs.getInt(InstPerformanceTable.COLUMN_COST_TOTAL); + return new InstPerformance(instanceId, callTimes, costTotal); + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } return null; } @Override public int getTpsMetric(int instanceId, long timeBucket) { + H2Client client = getClient(); + String sql = MessageFormat.format(GET_TPS_METRIC_SQL, InstPerformanceTable.TABLE, "id"); + Object[] params = new Object[]{instanceId}; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + return rs.getInt(InstPerformanceTable.COLUMN_CALLS); + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } return 0; } @Override public JsonArray getTpsMetric(int instanceId, long startTimeBucket, long endTimeBucket) { - return null; + H2Client client = getClient(); + String sql = MessageFormat.format(GET_TPS_METRICS_SQL, InstPerformanceTable.TABLE, "id"); + + long timeBucket = startTimeBucket; + List idList = new ArrayList<>(); + do { + String id = timeBucket + Const.ID_SPLIT + instanceId; + timeBucket = TimeBucketUtils.INSTANCE.addSecondForSecondTimeBucket(TimeBucketUtils.TimeBucketType.SECOND.name(), timeBucket, 1); + idList.add(id); + } + while (timeBucket <= endTimeBucket); + + StringBuilder builder = new StringBuilder(); + for( int i = 0 ; i < idList.size(); i++ ) { + builder.append("?,"); + } + builder.delete(builder.length() - 1, builder.length()); + builder.append(")"); + sql = sql + builder; + Object[] params = idList.toArray(new String[0]); + + JsonArray metrics = new JsonArray(); + try (ResultSet rs = client.executeQuery(sql, params)) { + while (rs.next()) { + int calls = rs.getInt(InstPerformanceTable.COLUMN_CALLS); + metrics.add(calls); + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } + return metrics; } @Override public int getRespTimeMetric(int instanceId, long timeBucket) { + H2Client client = getClient(); + String sql = MessageFormat.format(GET_TPS_METRIC_SQL, InstPerformanceTable.TABLE, "id"); + Object[] params = new Object[]{instanceId}; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + int callTimes = rs.getInt(InstPerformanceTable.COLUMN_CALLS); + int costTotal = rs.getInt(InstPerformanceTable.COLUMN_COST_TOTAL); + return costTotal / callTimes; + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } return 0; } @Override public JsonArray getRespTimeMetric(int instanceId, long startTimeBucket, long endTimeBucket) { - return null; + H2Client client = getClient(); + String sql = MessageFormat.format(GET_TPS_METRICS_SQL, InstPerformanceTable.TABLE, "id"); + + long timeBucket = startTimeBucket; + List idList = new ArrayList<>(); + do { + String id = timeBucket + Const.ID_SPLIT + instanceId; + timeBucket = TimeBucketUtils.INSTANCE.addSecondForSecondTimeBucket(TimeBucketUtils.TimeBucketType.SECOND.name(), timeBucket, 1); + idList.add(id); + } + while (timeBucket <= endTimeBucket); + + StringBuilder builder = new StringBuilder(); + for( int i = 0 ; i < idList.size(); i++ ) { + builder.append("?,"); + } + builder.delete(builder.length() - 1, builder.length()); + builder.append(")"); + sql = sql + builder; + Object[] params = idList.toArray(new String[0]); + + JsonArray metrics = new JsonArray(); + try (ResultSet rs = client.executeQuery(sql, params)) { + while (rs.next()) { + int callTimes = rs.getInt(InstPerformanceTable.COLUMN_CALLS); + int costTotal = rs.getInt(InstPerformanceTable.COLUMN_COST_TOTAL); + metrics.add(costTotal / callTimes); + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } + return metrics; } } 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 0f7028bbe..02412a367 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,31 +1,135 @@ package org.skywalking.apm.collector.ui.dao; import com.google.gson.JsonArray; + +import java.sql.ResultSet; +import java.sql.SQLException; +import java.text.MessageFormat; +import java.util.LinkedList; import java.util.List; + +import com.google.gson.JsonObject; +import org.skywalking.apm.collector.client.h2.H2Client; +import org.skywalking.apm.collector.client.h2.H2ClientException; +import org.skywalking.apm.collector.core.util.StringUtils; +import org.skywalking.apm.collector.core.util.TimeBucketUtils; +import org.skywalking.apm.collector.storage.define.node.NodeComponentTable; import org.skywalking.apm.collector.storage.define.register.InstanceDataDefine; +import org.skywalking.apm.collector.storage.define.register.InstanceTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.ui.cache.ApplicationCache; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author pengys5 */ public class InstanceH2DAO extends H2DAO implements IInstanceDAO { - @Override public Long lastHeartBeatTime() { + private final Logger logger = LoggerFactory.getLogger(InstanceH2DAO.class); + private static final String GET_LAST_HEARTBEAT_TIME_SQL = "select {0} from {1} where {2} > ? limit 1"; + private static final String GET_INST_LAST_HEARTBEAT_TIME_SQL = "select {0} from {1} where {2} > ? and {3} = ? limit 1"; + private static final String GET_INSTANCE_SQL = "select * from {0} where {1} = ?"; + private static final String GET_INSTANCES_SQL = "select * from {0} where {1} = ? and {2} >= ?"; + private static final String GET_APPLICATIONS_SQL = "select {3}, count({0}) as cnt from {1} where {2} >= ? and {2} <= ? group by {3} limit 100"; + + @Override + public Long lastHeartBeatTime() { + H2Client client = getClient(); + long fiveMinuteBefore = System.currentTimeMillis() - 5 * 60 * 1000; + fiveMinuteBefore = TimeBucketUtils.INSTANCE.getSecondTimeBucket(fiveMinuteBefore); + String sql = MessageFormat.format(GET_LAST_HEARTBEAT_TIME_SQL, InstanceTable.COLUMN_HEARTBEAT_TIME, InstanceTable.TABLE, InstanceTable.COLUMN_HEARTBEAT_TIME); + Object[] params = new Object[]{fiveMinuteBefore}; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + return rs.getLong(1); + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } + return 0L; + } + + @Override + public Long instanceLastHeartBeatTime(long applicationInstanceId) { + H2Client client = getClient(); + long fiveMinuteBefore = System.currentTimeMillis() - 5 * 60 * 1000; + fiveMinuteBefore = TimeBucketUtils.INSTANCE.getSecondTimeBucket(fiveMinuteBefore); + String sql = MessageFormat.format(GET_INST_LAST_HEARTBEAT_TIME_SQL, InstanceTable.COLUMN_HEARTBEAT_TIME, InstanceTable.TABLE, + InstanceTable.COLUMN_HEARTBEAT_TIME, InstanceTable.COLUMN_APPLICATION_ID); + Object[] params = new Object[]{fiveMinuteBefore, applicationInstanceId}; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + return rs.getLong(1); + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } + return 0L; + } + + @Override + public JsonArray getApplications(long startTime, long endTime) { + H2Client client = getClient(); + JsonArray applications = new JsonArray(); + String sql = MessageFormat.format(GET_APPLICATIONS_SQL, InstanceTable.COLUMN_INSTANCE_ID, + InstanceTable.TABLE, InstanceTable.COLUMN_HEARTBEAT_TIME, InstanceTable.COLUMN_APPLICATION_ID); + Object[] params = new Object[]{startTime, endTime}; + try (ResultSet rs = client.executeQuery(sql, params)) { + while (rs.next()) { + Integer applicationId = rs.getInt(InstanceTable.COLUMN_APPLICATION_ID); + logger.debug("applicationId: {}", applicationId); + JsonObject application = new JsonObject(); + application.addProperty("applicationId", applicationId); + application.addProperty("applicationCode", ApplicationCache.getForUI(applicationId)); + application.addProperty("instanceCount", rs.getInt("cnt")); + applications.add(application); + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } + return applications; + } + + @Override + public InstanceDataDefine.Instance getInstance(int instanceId) { + H2Client client = getClient(); + String sql = MessageFormat.format(GET_INSTANCE_SQL, InstanceTable.TABLE, InstanceTable.COLUMN_INSTANCE_ID); + Object[] params = new Object[]{instanceId}; + try (ResultSet rs = client.executeQuery(sql, params)) { + if (rs.next()) { + InstanceDataDefine.Instance instance = new InstanceDataDefine.Instance(); + instance.setId(String.valueOf(instanceId)); + instance.setApplicationId(rs.getInt(InstanceTable.COLUMN_APPLICATION_ID)); + instance.setAgentUUID(rs.getString(InstanceTable.COLUMN_AGENT_UUID)); + instance.setRegisterTime(rs.getLong(InstanceTable.COLUMN_REGISTER_TIME)); + instance.setHeartBeatTime(rs.getLong(InstanceTable.COLUMN_HEARTBEAT_TIME)); + instance.setOsInfo(rs.getString(InstanceTable.COLUMN_OS_INFO)); + return instance; + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } return null; } - @Override public Long instanceLastHeartBeatTime(long applicationInstanceId) { - return null; - } - - @Override public JsonArray getApplications(long startTime, long endTime) { - return null; - } - - @Override public InstanceDataDefine.Instance getInstance(int instanceId) { - return null; - } - - @Override public List getInstances(int applicationId, long timeBucket) { - return null; + @Override + public List getInstances(int applicationId, long timeBucket) { + logger.debug("get instances info, application id: {}, timeBucket: {}", applicationId, timeBucket); + List instanceList = new LinkedList<>(); + H2Client client = getClient(); + String sql = MessageFormat.format(GET_INSTANCES_SQL, InstanceTable.TABLE, InstanceTable.COLUMN_APPLICATION_ID, InstanceTable.COLUMN_HEARTBEAT_TIME); + Object[] params = new Object[]{applicationId, timeBucket}; + try (ResultSet rs = client.executeQuery(sql, params)) { + while (rs.next()) { + InstanceDataDefine.Instance instance = new InstanceDataDefine.Instance(); + instance.setApplicationId(rs.getInt(InstanceTable.COLUMN_APPLICATION_ID)); + instance.setHeartBeatTime(rs.getLong(InstanceTable.COLUMN_HEARTBEAT_TIME)); + instance.setInstanceId(rs.getInt(InstanceTable.COLUMN_INSTANCE_ID)); + instanceList.add(instance); + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } + return instanceList; } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeComponentH2DAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeComponentH2DAO.java index 5a73d42ce..e3f818691 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeComponentH2DAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeComponentH2DAO.java @@ -1,14 +1,62 @@ package org.skywalking.apm.collector.ui.dao; import com.google.gson.JsonArray; +import com.google.gson.JsonObject; +import org.skywalking.apm.collector.client.h2.H2Client; +import org.skywalking.apm.collector.client.h2.H2ClientException; +import org.skywalking.apm.collector.core.util.StringUtils; +import org.skywalking.apm.collector.storage.define.node.NodeComponentTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.ui.cache.ApplicationCache; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.sql.ResultSet; +import java.sql.SQLException; +import java.text.MessageFormat; /** * @author pengys5 */ public class NodeComponentH2DAO extends H2DAO implements INodeComponentDAO { - + private final Logger logger = LoggerFactory.getLogger(NodeComponentH2DAO.class); + private static final String AGGREGATE_COMPONENT_SQL = "select * from {3} where {4} >= ? and {4} <= ? group by {0}, {1}, {2} limit 100"; @Override public JsonArray load(long startTime, long endTime) { - return null; + JsonArray nodeComponentArray = new JsonArray(); + nodeComponentArray.addAll(aggregationComponent(startTime, endTime)); + return nodeComponentArray; + } + private JsonArray aggregationComponent(long startTime, long endTime) { + H2Client client = getClient(); + + JsonArray nodeComponentArray = new JsonArray(); + String sql = MessageFormat.format(AGGREGATE_COMPONENT_SQL, NodeComponentTable.COLUMN_COMPONENT_ID, + NodeComponentTable.COLUMN_PEER, NodeComponentTable.COLUMN_PEER_ID, + NodeComponentTable.TABLE, NodeComponentTable.COLUMN_TIME_BUCKET); + Object[] params = new Object[]{startTime, endTime}; + try (ResultSet rs = client.executeQuery(sql, params)) { + while (rs.next()) { + int peerId = rs.getInt(NodeComponentTable.COLUMN_PEER_ID); + String componentName = rs.getString(NodeComponentTable.COLUMN_COMPONENT_NAME); + if (peerId != 0) { + String peer = ApplicationCache.getForUI(peerId); + nodeComponentArray.add(buildNodeComponent(peer, componentName)); + } + String peer = rs.getString(NodeComponentTable.COLUMN_PEER); + if (StringUtils.isNotEmpty(peer)) { + nodeComponentArray.add(buildNodeComponent(peer, componentName)); + } + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } + return nodeComponentArray; + } + + private JsonObject buildNodeComponent(String peer, String componentName) { + JsonObject nodeComponentObj = new JsonObject(); + nodeComponentObj.addProperty("componentName", componentName); + nodeComponentObj.addProperty("peer", peer); + return nodeComponentObj; } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeMappingH2DAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeMappingH2DAO.java index afda33070..6fb2cea3c 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeMappingH2DAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeMappingH2DAO.java @@ -1,14 +1,58 @@ package org.skywalking.apm.collector.ui.dao; import com.google.gson.JsonArray; +import com.google.gson.JsonObject; +import org.skywalking.apm.collector.client.h2.H2Client; +import org.skywalking.apm.collector.client.h2.H2ClientException; +import org.skywalking.apm.collector.core.util.StringUtils; +import org.skywalking.apm.collector.storage.define.node.NodeMappingTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.ui.cache.ApplicationCache; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.sql.ResultSet; +import java.sql.SQLException; +import java.text.MessageFormat; /** * @author pengys5 */ public class NodeMappingH2DAO extends H2DAO implements INodeMappingDAO { - + private final Logger logger = LoggerFactory.getLogger(NodeMappingH2DAO.class); + private static final String NODE_MAPPING_SQL = "select * from {3} where {4} >= ? and {4} <= ? group by {0}, {1}, {2} limit 100"; @Override public JsonArray load(long startTime, long endTime) { - return null; + H2Client client = getClient(); + JsonArray nodeMappingArray = new JsonArray(); + String sql = MessageFormat.format(NODE_MAPPING_SQL, NodeMappingTable.COLUMN_APPLICATION_ID, + NodeMappingTable.COLUMN_ADDRESS_ID, NodeMappingTable.COLUMN_ADDRESS, + NodeMappingTable.TABLE, NodeMappingTable.COLUMN_TIME_BUCKET); + + Object[] params = new Object[]{startTime, endTime}; + try (ResultSet rs = client.executeQuery(sql, params)) { + while (rs.next()) { + int applicationId = rs.getInt(NodeMappingTable.COLUMN_APPLICATION_ID); + String applicationCode = ApplicationCache.getForUI(applicationId); + int addressId = rs.getInt(NodeMappingTable.COLUMN_ADDRESS_ID); + if (addressId != 0) { + String address = ApplicationCache.getForUI(addressId); + JsonObject nodeMappingObj = new JsonObject(); + nodeMappingObj.addProperty("applicationCode", applicationCode); + nodeMappingObj.addProperty("address", address); + nodeMappingArray.add(nodeMappingObj); + } + String address = rs.getString(NodeMappingTable.COLUMN_ADDRESS); + if (StringUtils.isNotEmpty(address)) { + JsonObject nodeMappingObj = new JsonObject(); + nodeMappingObj.addProperty("applicationCode", applicationCode); + nodeMappingObj.addProperty("address", address); + nodeMappingArray.add(nodeMappingObj); + } + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } + logger.debug("node mapping data: {}", nodeMappingArray.toString()); + return nodeMappingArray; } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeReferenceH2DAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeReferenceH2DAO.java index 72c6cd65e..ab9df2fdb 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeReferenceH2DAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/NodeReferenceH2DAO.java @@ -1,13 +1,74 @@ package org.skywalking.apm.collector.ui.dao; import com.google.gson.JsonArray; +import com.google.gson.JsonObject; +import org.skywalking.apm.collector.client.h2.H2Client; +import org.skywalking.apm.collector.client.h2.H2ClientException; +import org.skywalking.apm.collector.core.util.StringUtils; +import org.skywalking.apm.collector.storage.define.node.NodeMappingTable; +import org.skywalking.apm.collector.storage.define.noderef.NodeReferenceTable; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; +import org.skywalking.apm.collector.ui.cache.ApplicationCache; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.sql.ResultSet; +import java.sql.SQLException; +import java.text.MessageFormat; /** * @author pengys5 */ public class NodeReferenceH2DAO extends H2DAO implements INodeReferenceDAO { + private final Logger logger = LoggerFactory.getLogger(NodeReferenceH2DAO.class); + private static final String NODE_REFERENCE_SQL = "select {8}, {9}, {10}, sum({0}) as {0}, sum({1}) as {1}, sum({2}) as {2}, " + + "sum({3}) as {3}, sum({4}) as {4}, sum({5}) as {5} from {6} where {7} >= ? and {7} <= ? group by {8}, {9}, {10} limit 100"; @Override public JsonArray load(long startTime, long endTime) { - return null; + H2Client client = getClient(); + JsonArray nodeRefResSumArray = new JsonArray(); + String sql = MessageFormat.format(NODE_REFERENCE_SQL, NodeReferenceTable.COLUMN_S1_LTE, + NodeReferenceTable.COLUMN_S3_LTE, NodeReferenceTable.COLUMN_S5_LTE, + NodeReferenceTable.COLUMN_S5_GT, NodeReferenceTable.COLUMN_SUMMARY, + NodeReferenceTable.COLUMN_ERROR, NodeReferenceTable.TABLE, NodeReferenceTable.COLUMN_TIME_BUCKET, + NodeReferenceTable.COLUMN_FRONT_APPLICATION_ID, NodeReferenceTable.COLUMN_BEHIND_APPLICATION_ID, NodeReferenceTable.COLUMN_BEHIND_PEER); + + Object[] params = new Object[]{startTime, endTime}; + try (ResultSet rs = client.executeQuery(sql, params)) { + while (rs.next()) { + int applicationId = rs.getInt(NodeReferenceTable.COLUMN_FRONT_APPLICATION_ID); + String applicationCode = ApplicationCache.getForUI(applicationId); + int behindApplicationId = rs.getInt(NodeReferenceTable.COLUMN_BEHIND_APPLICATION_ID); + if (behindApplicationId != 0) { + String behindApplicationCode = ApplicationCache.getForUI(behindApplicationId); + + JsonObject nodeRefResSumObj = new JsonObject(); + nodeRefResSumObj.addProperty("front", applicationCode); + nodeRefResSumObj.addProperty("behind", behindApplicationCode); + nodeRefResSumObj.addProperty(NodeReferenceTable.COLUMN_S1_LTE, rs.getDouble(NodeReferenceTable.COLUMN_S1_LTE)); + nodeRefResSumObj.addProperty(NodeReferenceTable.COLUMN_S3_LTE, rs.getDouble(NodeReferenceTable.COLUMN_S3_LTE)); + nodeRefResSumObj.addProperty(NodeReferenceTable.COLUMN_S5_LTE, rs.getDouble(NodeReferenceTable.COLUMN_S5_LTE)); + nodeRefResSumObj.addProperty(NodeReferenceTable.COLUMN_S5_GT, rs.getDouble(NodeReferenceTable.COLUMN_S5_GT)); + nodeRefResSumObj.addProperty(NodeReferenceTable.COLUMN_ERROR, rs.getDouble(NodeReferenceTable.COLUMN_ERROR)); + nodeRefResSumObj.addProperty(NodeReferenceTable.COLUMN_SUMMARY, rs.getDouble(NodeReferenceTable.COLUMN_SUMMARY)); + nodeRefResSumArray.add(nodeRefResSumObj); + } + String behindPeer = rs.getString(NodeReferenceTable.COLUMN_BEHIND_PEER); + if (StringUtils.isNotEmpty(behindPeer)) { + JsonObject nodeRefResSumObj = new JsonObject(); + nodeRefResSumObj.addProperty("front", applicationCode); + nodeRefResSumObj.addProperty("behind", behindPeer); + nodeRefResSumObj.addProperty(NodeReferenceTable.COLUMN_S1_LTE, rs.getDouble(NodeReferenceTable.COLUMN_S1_LTE)); + nodeRefResSumObj.addProperty(NodeReferenceTable.COLUMN_S3_LTE, rs.getDouble(NodeReferenceTable.COLUMN_S3_LTE)); + nodeRefResSumObj.addProperty(NodeReferenceTable.COLUMN_S5_LTE, rs.getDouble(NodeReferenceTable.COLUMN_S5_LTE)); + nodeRefResSumObj.addProperty(NodeReferenceTable.COLUMN_S5_GT, rs.getDouble(NodeReferenceTable.COLUMN_S5_GT)); + nodeRefResSumObj.addProperty(NodeReferenceTable.COLUMN_ERROR, rs.getDouble(NodeReferenceTable.COLUMN_ERROR)); + nodeRefResSumObj.addProperty(NodeReferenceTable.COLUMN_SUMMARY, rs.getDouble(NodeReferenceTable.COLUMN_SUMMARY)); + nodeRefResSumArray.add(nodeRefResSumObj); + } + } + } catch (SQLException | H2ClientException e) { + logger.error(e.getMessage(), e); + } + return nodeRefResSumArray; } } 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 a29e509cc..1d5a0ee0b 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 @@ -41,10 +41,14 @@ public class InstanceHealthService { IGCMetricDAO gcMetricDAO = (IGCMetricDAO)DAOContainer.INSTANCE.get(IGCMetricDAO.class.getName()); JsonObject instanceJson = new JsonObject(); instanceJson.addProperty("id", instance.getInstanceId()); - instanceJson.addProperty("tps", performance.getCalls()); + if (performance != null) { + instanceJson.addProperty("tps", performance.getCalls()); + } else { + instanceJson.addProperty("tps", 0); + } int avg = 0; - if (performance.getCalls() != 0) { + if (performance != null && performance.getCalls() != 0) { avg = (int)(performance.getCostTotal() / performance.getCalls()); } instanceJson.addProperty("avg", avg);