add h2 support ga
This commit is contained in:
parent
5219e23d52
commit
bc38a56da0
|
|
@ -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<String, Object>, Map<String, Object>> {
|
||||
public class CpuMetricH2DAO extends H2DAO implements ICpuMetricDAO, IPersistenceDAO<H2SqlEntity, H2SqlEntity> {
|
||||
private final Logger logger = LoggerFactory.getLogger(CpuMetricH2DAO.class);
|
||||
@Override public Data get(String id, DataDefine dataDefine) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override public Map<String, Object> prepareBatchInsert(Data data) {
|
||||
@Override public H2SqlEntity prepareBatchInsert(Data data) {
|
||||
H2SqlEntity entity = new H2SqlEntity();
|
||||
Map<String, Object> 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<String, Object> prepareBatchUpdate(Data data) {
|
||||
@Override public H2SqlEntity prepareBatchUpdate(Data data) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object>, Map<String, Object>> {
|
||||
public class GCMetricH2DAO extends H2DAO implements IGCMetricDAO, IPersistenceDAO<H2SqlEntity, H2SqlEntity> {
|
||||
@Override public Data get(String id, DataDefine dataDefine) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override public Map<String, Object> prepareBatchInsert(Data data) {
|
||||
@Override public H2SqlEntity prepareBatchInsert(Data data) {
|
||||
H2SqlEntity entity = new H2SqlEntity();
|
||||
Map<String, Object> 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<String, Object> prepareBatchUpdate(Data data) {
|
||||
@Override public H2SqlEntity prepareBatchUpdate(Data data) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object>, Map<String, Object>> {
|
||||
public class InstanceHeartBeatH2DAO extends H2DAO implements IInstanceHeartBeatDAO, IPersistenceDAO<H2SqlEntity, H2SqlEntity> {
|
||||
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<String, Object> 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<String, Object> prepareBatchUpdate(Data data) {
|
||||
@Override public H2SqlEntity prepareBatchUpdate(Data data) {
|
||||
H2SqlEntity entity = new H2SqlEntity();
|
||||
Map<String, Object> 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<Object> params = new ArrayList<>(source.values());
|
||||
params.add(data.getDataString(0));
|
||||
entity.setParams(params.toArray(new Object[0]));
|
||||
return entity;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object>, Map<String, Object>> {
|
||||
public class MemoryMetricH2DAO extends H2DAO implements IMemoryMetricDAO, IPersistenceDAO<H2SqlEntity, H2SqlEntity> {
|
||||
@Override public Data get(String id, DataDefine dataDefine) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override public Map<String, Object> prepareBatchInsert(Data data) {
|
||||
@Override public H2SqlEntity prepareBatchInsert(Data data) {
|
||||
H2SqlEntity entity = new H2SqlEntity();
|
||||
Map<String, Object> 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<String, Object> prepareBatchUpdate(Data data) {
|
||||
@Override public H2SqlEntity prepareBatchUpdate(Data data) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object>, Map<String, Object>> {
|
||||
public class MemoryPoolMetricH2DAO extends H2DAO implements IMemoryPoolMetricDAO, IPersistenceDAO<H2SqlEntity, H2SqlEntity> {
|
||||
@Override public Data get(String id, DataDefine dataDefine) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override public Map<String, Object> prepareBatchInsert(Data data) {
|
||||
@Override public H2SqlEntity prepareBatchInsert(Data data) {
|
||||
H2SqlEntity entity = new H2SqlEntity();
|
||||
Map<String, Object> 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<String, Object> prepareBatchUpdate(Data data) {
|
||||
@Override public H2SqlEntity prepareBatchUpdate(Data data) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object> 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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object>, Map<String, Object>> {
|
||||
public class GlobalTraceH2DAO extends H2DAO implements IGlobalTraceDAO, IPersistenceDAO<H2SqlEntity, H2SqlEntity> {
|
||||
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<String, Object> 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<String, Object> prepareBatchUpdate(Data data) {
|
||||
return null;
|
||||
|
||||
@Override public H2SqlEntity prepareBatchInsert(Data data) {
|
||||
Map<String, Object> 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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object>, Map<String, Object>> {
|
||||
public class InstPerformanceH2DAO extends H2DAO implements IInstPerformanceDAO, IPersistenceDAO<H2SqlEntity, H2SqlEntity> {
|
||||
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<String, Object> prepareBatchInsert(Data data) {
|
||||
return null;
|
||||
@Override public H2SqlEntity prepareBatchInsert(Data data) {
|
||||
Map<String, Object> 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<String, Object> prepareBatchUpdate(Data data) {
|
||||
return null;
|
||||
@Override public H2SqlEntity prepareBatchUpdate(Data data) {
|
||||
Map<String, Object> 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<Object> values = new ArrayList<>(source.values());
|
||||
values.add(id);
|
||||
entity.setParams(values.toArray(new Object[0]));
|
||||
return entity;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object>, Map<String, Object>> {
|
||||
public class NodeComponentH2DAO extends H2DAO implements INodeComponentDAO, IPersistenceDAO<H2SqlEntity, H2SqlEntity> {
|
||||
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<String, Object> prepareBatchInsert(Data data) {
|
||||
return null;
|
||||
|
||||
@Override
|
||||
public H2SqlEntity prepareBatchInsert(Data data) {
|
||||
Map<String, Object> 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<String, Object> prepareBatchUpdate(Data data) {
|
||||
return null;
|
||||
|
||||
@Override
|
||||
public H2SqlEntity prepareBatchUpdate(Data data) {
|
||||
Map<String, Object> 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<Object> values = new ArrayList<>(source.values());
|
||||
values.add(id);
|
||||
entity.setParams(values.toArray(new Object[0]));
|
||||
return entity;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object>, Map<String, Object>> {
|
||||
public class NodeMappingH2DAO extends H2DAO implements INodeMappingDAO, IPersistenceDAO<H2SqlEntity, H2SqlEntity> {
|
||||
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<String, Object> prepareBatchInsert(Data data) {
|
||||
return null;
|
||||
@Override public H2SqlEntity prepareBatchInsert(Data data) {
|
||||
Map<String, Object> 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<String, Object> prepareBatchUpdate(Data data) {
|
||||
return null;
|
||||
@Override public H2SqlEntity prepareBatchUpdate(Data data) {
|
||||
Map<String, Object> 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<Object> values = new ArrayList<>(source.values());
|
||||
values.add(id);
|
||||
entity.setParams(values.toArray(new Object[0]));
|
||||
return entity;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object>, Map<String, Object>> {
|
||||
public class NodeReferenceH2DAO extends H2DAO implements INodeReferenceDAO, IPersistenceDAO<H2SqlEntity, H2SqlEntity> {
|
||||
private final Logger logger = LoggerFactory.getLogger(NodeReferenceH2DAO.class);
|
||||
@Override public Data get(String id, DataDefine dataDefine) {
|
||||
return null;
|
||||
}
|
||||
@Override public Map<String, Object> prepareBatchInsert(Data data) {
|
||||
@Override public H2SqlEntity prepareBatchInsert(Data data) {
|
||||
return null;
|
||||
}
|
||||
@Override public Map<String, Object> prepareBatchUpdate(Data data) {
|
||||
@Override public H2SqlEntity prepareBatchUpdate(Data data) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object>, Map<String, Object>> {
|
||||
public class SegmentCostH2DAO extends H2DAO implements ISegmentCostDAO, IPersistenceDAO<H2SqlEntity, H2SqlEntity> {
|
||||
private final Logger logger = LoggerFactory.getLogger(SegmentCostH2DAO.class);
|
||||
@Override public Data get(String id, DataDefine dataDefine) {
|
||||
return null;
|
||||
}
|
||||
@Override public Map<String, Object> 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<String, Object> 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<String, Object> prepareBatchUpdate(Data data) {
|
||||
@Override public H2SqlEntity prepareBatchUpdate(Data data) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object>, Map<String, Object>> {
|
||||
public class SegmentH2DAO extends H2DAO implements ISegmentDAO, IPersistenceDAO<H2SqlEntity, H2SqlEntity> {
|
||||
private final Logger logger = LoggerFactory.getLogger(SegmentCostH2DAO.class);
|
||||
@Override public Data get(String id, DataDefine dataDefine) {
|
||||
return null;
|
||||
}
|
||||
@Override public Map<String, Object> prepareBatchInsert(Data data) {
|
||||
return null;
|
||||
@Override public H2SqlEntity prepareBatchInsert(Data data) {
|
||||
Map<String, Object> 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<String, Object> prepareBatchUpdate(Data data) {
|
||||
@Override public H2SqlEntity prepareBatchUpdate(Data data) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object>, Map<String, Object>> {
|
||||
public class ServiceEntryH2DAO extends H2DAO implements IServiceEntryDAO, IPersistenceDAO<H2SqlEntity, H2SqlEntity> {
|
||||
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<String, Object> prepareBatchInsert(Data data) {
|
||||
@Override public H2SqlEntity prepareBatchInsert(Data data) {
|
||||
H2SqlEntity entity = new H2SqlEntity();
|
||||
Map<String, Object> 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<String, Object> prepareBatchUpdate(Data data) {
|
||||
@Override public H2SqlEntity prepareBatchUpdate(Data data) {
|
||||
H2SqlEntity entity = new H2SqlEntity();
|
||||
Map<String, Object> 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<Object> values = new ArrayList<>(source.values());
|
||||
values.add(id);
|
||||
entity.setParams(values.toArray(new Object[0]));
|
||||
return entity;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object>, Map<String, Object>> {
|
||||
public class ServiceReferenceH2DAO extends H2DAO implements IServiceReferenceDAO, IPersistenceDAO<H2SqlEntity, H2SqlEntity> {
|
||||
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<String, Object> prepareBatchInsert(Data data) {
|
||||
@Override
|
||||
public H2SqlEntity prepareBatchInsert(Data data) {
|
||||
H2SqlEntity entity = new H2SqlEntity();
|
||||
Map<String, Object> 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<String, Object> prepareBatchUpdate(Data data) {
|
||||
@Override
|
||||
public H2SqlEntity prepareBatchUpdate(Data data) {
|
||||
H2SqlEntity entity = new H2SqlEntity();
|
||||
Map<String, Object> 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<Object> values = new ArrayList<>(source.values());
|
||||
values.add(id);
|
||||
entity.setParams(values.toArray(new Object[0]));
|
||||
return entity;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, PreparedStatement> 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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<H2Client> {
|
|||
|
||||
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<H2Client> {
|
|||
}
|
||||
} catch (SQLException | H2ClientException e) {
|
||||
logger.error(e.getMessage(), e);
|
||||
} finally {
|
||||
client.closeResultSet(rs);
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
public final String getBatchInsertSql(String tableName, Set<String> 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<String> 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();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
@ -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
|
||||
*/
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String> 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;
|
||||
}
|
||||
}
|
||||
|
|
@ -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) {
|
||||
|
|
|
|||
|
|
@ -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<String> 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<String> 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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<InstanceDataDefine.Instance> getInstances(int applicationId, long timeBucket) {
|
||||
return null;
|
||||
@Override
|
||||
public List<InstanceDataDefine.Instance> getInstances(int applicationId, long timeBucket) {
|
||||
logger.debug("get instances info, application id: {}, timeBucket: {}", applicationId, timeBucket);
|
||||
List<InstanceDataDefine.Instance> 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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
Loading…
Reference in New Issue