diff --git a/src/main/java/com/nhn/hippo/web/calltree/server/Server.java b/src/main/java/com/nhn/hippo/web/calltree/server/Server.java index 01b6ac1f2..19ef986e9 100644 --- a/src/main/java/com/nhn/hippo/web/calltree/server/Server.java +++ b/src/main/java/com/nhn/hippo/web/calltree/server/Server.java @@ -7,16 +7,18 @@ public class Server implements Comparable { private int sequence; private final String id; private final String agentId; + private final String applicationName; private final String endPoint; private final boolean terminal; - public Server(String agentId, String endPoint, boolean terminal) { -// this.id = agentId + ":" + endPoint; - this.id = endPoint; - this.agentId = agentId; - this.endPoint = endPoint; - this.terminal = terminal; - } + public Server(String agentId, String applicationName, String endPoint, boolean terminal) { + // this.id = agentId + ":" + endPoint; + this.id = endPoint; + this.agentId = agentId; + this.applicationName = applicationName; + this.endPoint = endPoint; + this.terminal = terminal; + } public String getId() { return this.id; @@ -42,7 +44,11 @@ public class Server implements Comparable { return terminal; } - @Override + public String getApplicationName() { + return applicationName; + } + + @Override public int compareTo(Server server) { return id.compareTo(server.id); } diff --git a/src/main/java/com/nhn/hippo/web/calltree/server/ServerCallTree.java b/src/main/java/com/nhn/hippo/web/calltree/server/ServerCallTree.java index bb4f2b59e..7e14cdcaf 100644 --- a/src/main/java/com/nhn/hippo/web/calltree/server/ServerCallTree.java +++ b/src/main/java/com/nhn/hippo/web/calltree/server/ServerCallTree.java @@ -36,7 +36,7 @@ public class ServerCallTree { * make Servers */ // TODO: 여기에서 이러지말고 수집할 때 처음부터 table에 저장해둘 수 있나?? - Server server = new Server(span.getAgentId(), span.getEndPoint(), span.isTerminal()); + Server server = new Server(span.getAgentId(), span.getApplicationName(), span.getEndPoint(), span.isTerminal()); if (server.getId() == null) { return; diff --git a/src/main/java/com/nhn/hippo/web/controller/FlowChartController.java b/src/main/java/com/nhn/hippo/web/controller/FlowChartController.java index 7b2a7b804..c627b1967 100644 --- a/src/main/java/com/nhn/hippo/web/controller/FlowChartController.java +++ b/src/main/java/com/nhn/hippo/web/controller/FlowChartController.java @@ -32,8 +32,7 @@ public class FlowChartController { @RequestMapping(value = "/flow", method = RequestMethod.GET) public String flow(Model model, @RequestParam("application") String applicationName, @RequestParam("from") long from, @RequestParam("to") long to) { - String[] agentIds = flow.selectAgentIdsFromApplicationName(applicationName); - Set traceIds = flow.selectTraceIdsFromTraceIndex(agentIds, from, to); + Set traceIds = flow.selectTraceIdsFromApplicationTraceIndex(applicationName, from, to); RPCCallTree callTree = flow.selectRPCCallTree(traceIds); @@ -47,15 +46,62 @@ public class FlowChartController { @RequestMapping(value = "/flowserver", method = RequestMethod.GET) public String flowserver(Model model, @RequestParam("application") String applicationName, @RequestParam("from") long from, @RequestParam("to") long to) { - String[] agentIds = flow.selectAgentIdsFromApplicationName(applicationName); // TODO 제거 하거나, interceptor로 할것. StopWatch watch = new StopWatch(); watch.start("scanTraceindex"); - Set traceIds = flow.selectTraceIdsFromTraceIndex(agentIds, from, to); + + Set traceIds = flow.selectTraceIdsFromApplicationTraceIndex(applicationName, from, to); + watch.stop(); logger.info("time:{} {}", watch.getLastTaskTimeMillis(), traceIds.size()); watch.start("selectServerCallTree"); + ServerCallTree callTree = flow.selectServerCallTree(traceIds); + + watch.stop(); + logger.info("time:{}", watch.getLastTaskTimeMillis()); + + model.addAttribute("nodes", callTree.getNodes()); + model.addAttribute("links", callTree.getLinks()); + model.addAttribute("businessTransactions", callTree.getBusinessTransactions().getBusinessTransactionIterator()); + model.addAttribute("traces", callTree.getBusinessTransactions().getTracesIterator()); + + logger.debug("callTree:{}", callTree); + + return "flowserver"; + } + + @RequestMapping(value = "/flow2", method = RequestMethod.GET) + public String flowbyHost(Model model, @RequestParam("host") String[] hosts, @RequestParam("from") long from, @RequestParam("to") long to) { + String[] agentIds = flow.selectAgentIds(hosts); + Set traceIds = flow.selectTraceIdsFromTraceIndex(agentIds, from, to); + + RPCCallTree callTree = flow.selectRPCCallTree(traceIds); + + model.addAttribute("nodes", callTree.getNodes()); + model.addAttribute("links", callTree.getLinks()); + + logger.debug("callTree:{}", callTree); + + return "flow"; + } + + @RequestMapping(value = "/flowserver2", method = RequestMethod.GET) + public String flowserverByHost(Model model, @RequestParam("host") String[] hosts, @RequestParam("from") long from, @RequestParam("to") long to) { + String[] agentIds = flow.selectAgentIds(hosts); + + // TODO 제거 하거나, interceptor로 할것. + StopWatch watch = new StopWatch(); + watch.start("scanTraceindex"); + + Set traceIds = flow.selectTraceIdsFromTraceIndex(agentIds, from, to); + + watch.stop(); + logger.info("time:{} {}", watch.getLastTaskTimeMillis(), traceIds.size()); + watch.start("selectServerCallTree"); + + ServerCallTree callTree = flow.selectServerCallTree(traceIds); + watch.stop(); logger.info("time:{}", watch.getLastTaskTimeMillis()); diff --git a/src/main/java/com/nhn/hippo/web/dao/ApplicationTraceIndexDao.java b/src/main/java/com/nhn/hippo/web/dao/ApplicationTraceIndexDao.java new file mode 100644 index 000000000..02a478e8a --- /dev/null +++ b/src/main/java/com/nhn/hippo/web/dao/ApplicationTraceIndexDao.java @@ -0,0 +1,12 @@ +package com.nhn.hippo.web.dao; + +import java.util.List; + +/** + * + */ +public interface ApplicationTraceIndexDao { + List scanTraceIndex(String applicationName, long start, long end); + + List> multiScanTraceIndex(String[] applicationNames, long start, long end); +} diff --git a/src/main/java/com/nhn/hippo/web/dao/HbaseApplicationIndexDao.java b/src/main/java/com/nhn/hippo/web/dao/HbaseApplicationIndexDao.java deleted file mode 100644 index ef4adcd50..000000000 --- a/src/main/java/com/nhn/hippo/web/dao/HbaseApplicationIndexDao.java +++ /dev/null @@ -1,52 +0,0 @@ -package com.nhn.hippo.web.dao; - -import java.util.List; - -import org.apache.hadoop.hbase.client.Get; -import org.apache.hadoop.hbase.client.Scan; -import org.apache.hadoop.hbase.util.Bytes; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.beans.factory.annotation.Qualifier; -import org.springframework.data.hadoop.hbase.RowMapper; -import org.springframework.stereotype.Repository; - -import com.profiler.common.hbase.HBaseTables; -import com.profiler.common.hbase.HbaseOperations2; - -/** - * @author netspider - */ -@Repository -public class HbaseApplicationIndexDao implements ApplicationIndexDao { - private final Logger logger = LoggerFactory.getLogger(this.getClass().getName()); - - @Autowired - private HbaseOperations2 hbaseOperations2; - - @Autowired - @Qualifier("applicationNameMapper") - private RowMapper applicationNameMapper; - - @Autowired - @Qualifier("agentIdMapper") - private RowMapper agentIdMapper; - - @Override - public List selectAllApplicationNames() { - Scan scan = new Scan(); - scan.setCaching(30); - return hbaseOperations2.find(HBaseTables.APPLICATION_INDEX, scan, applicationNameMapper); - } - - @Override - public String[] selectAgentIds(String applicationName) { - byte[] rowKey = Bytes.toBytes(applicationName); - - Get get = new Get(rowKey); - get.addFamily(HBaseTables.APPLICATION_CF_AGENTS); - - return hbaseOperations2.get(HBaseTables.APPLICATION_INDEX, get, agentIdMapper); - } -} diff --git a/src/main/java/com/nhn/hippo/web/dao/HbaseRootTraceIndexDao.java b/src/main/java/com/nhn/hippo/web/dao/HbaseRootTraceIndexDao.java deleted file mode 100644 index 662b228b1..000000000 --- a/src/main/java/com/nhn/hippo/web/dao/HbaseRootTraceIndexDao.java +++ /dev/null @@ -1,84 +0,0 @@ -package com.nhn.hippo.web.dao; - -import com.profiler.common.hbase.HBaseTables; -import com.profiler.common.hbase.HbaseOperations2; -import com.profiler.common.util.SpanUtils; -import org.apache.hadoop.hbase.client.Scan; -import org.apache.hadoop.hbase.util.Bytes; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.beans.factory.annotation.Qualifier; -import org.springframework.data.hadoop.hbase.RowMapper; -import org.springframework.stereotype.Repository; - -import java.util.ArrayList; -import java.util.List; - -/** - * - */ -@Repository -public class HbaseRootTraceIndexDao implements RootTraceIndexDao { - - private Logger logger = LoggerFactory.getLogger(this.getClass()); - - private final byte[] COLFAM_TRACE = Bytes.toBytes("Trace"); - private final byte[] COLNAME_ID = Bytes.toBytes("ID"); - - @Autowired - private HbaseOperations2 hbaseOperations2; - - @Autowired - @Qualifier("traceIndexMapper") - private RowMapper traceIndexMapper; - - - private int scanCacheSize = 40; - - public void setScanCacheSize(int scanCacheSize) { - this.scanCacheSize = scanCacheSize; - } - - - @Override - public List scanTraceIndex(String agent, long start, long end) { - Scan scan = createScan(agent, start, end); - return hbaseOperations2.find(HBaseTables.ROOT_TRACE_INDEX, scan, traceIndexMapper); - } - - @Override - public List> multiScanTraceIndex(String[] agents, long start, long end) { - final List multiScan = new ArrayList(agents.length); - for (String agent : agents) { - Scan scan = createScan(agent, start, end); - multiScan.add(scan); - } - return hbaseOperations2.find(HBaseTables.ROOT_TRACE_INDEX, multiScan, traceIndexMapper); - } - - private Scan createScan(String agent, long start, long end) { - Scan scan = new Scan(); - // cache size를 지정해야 되는거 같음.?? - scan.setCaching(this.scanCacheSize); - - byte[] bAgent = Bytes.toBytes(agent); - byte[] traceIndexStartKey = SpanUtils.getTraceIndexRowKey(bAgent, start); - scan.setStartRow(traceIndexStartKey); - - byte[] traceIndexEndKey = SpanUtils.getTraceIndexRowKey(bAgent, end); - scan.setStopRow(traceIndexEndKey); - - scan.addColumn(COLFAM_TRACE, COLNAME_ID); - scan.setId("rootTraceIndexScan"); - - // json으로 변화해서 로그를 찍어서. 최초 변환 속도가 느림. - logger.debug("create scan:{}", scan); - return scan; - } - - @Override - public List parallelScanTraceIndex(String[] agents, long start, long end) { - return null; //To change body of implemented methods use File | Settings | File Templates. - } -} diff --git a/src/main/java/com/nhn/hippo/web/dao/HbaseTraceDao.java b/src/main/java/com/nhn/hippo/web/dao/HbaseTraceDao.java deleted file mode 100644 index 6b78e2039..000000000 --- a/src/main/java/com/nhn/hippo/web/dao/HbaseTraceDao.java +++ /dev/null @@ -1,98 +0,0 @@ -package com.nhn.hippo.web.dao; - -import java.util.ArrayList; -import java.util.List; -import java.util.Set; -import java.util.UUID; - -import org.apache.hadoop.hbase.client.Get; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.beans.factory.annotation.Qualifier; -import org.springframework.data.hadoop.hbase.RowMapper; -import org.springframework.stereotype.Repository; - -import com.nhn.hippo.web.vo.TraceId; -import com.profiler.common.bo.SpanBo; -import com.profiler.common.hbase.HBaseTables; -import com.profiler.common.hbase.HbaseOperations2; -import com.profiler.common.util.BytesUtils; - -/** - * - */ -@Repository -public class HbaseTraceDao implements TraceDao { - - private final byte[] COLFAM_SPAN = HBaseTables.TRACES_CF_SPAN; - - private final byte[] COLFAM_ANNOTATION = HBaseTables.TRACES_CF_ANNOTATION; - - private Logger logger = LoggerFactory.getLogger(this.getClass()); - - @Autowired - private HbaseOperations2 template2; - - @Autowired - @Qualifier("spanMapper") - private RowMapper> spanMapper; - - @Autowired - @Qualifier("spanAnnotationMapper") - private RowMapper> spanAnnotationMapper; - - @Override - public List selectSpan(UUID traceId) { - byte[] uuidBytes = BytesUtils.longLongToBytes(traceId.getMostSignificantBits(), traceId.getLeastSignificantBits()); - return template2.get(HBaseTables.TRACES, uuidBytes, COLFAM_SPAN, spanMapper); - } - - public List selectSpanAndAnnotation(UUID traceId) { - byte[] uuidBytes = BytesUtils.longLongToBytes(traceId.getMostSignificantBits(), traceId.getLeastSignificantBits()); - Get get = new Get(uuidBytes); - get.addFamily(COLFAM_SPAN); - get.addFamily(COLFAM_ANNOTATION); - return template2.get(HBaseTables.TRACES, get, spanAnnotationMapper); - } - - @Override - public List selectSpan(long traceIdMost, long traceIdLeast) { - byte[] uuidBytes = BytesUtils.longLongToBytes(traceIdMost, traceIdLeast); - return template2.get(HBaseTables.TRACES, uuidBytes, COLFAM_SPAN, spanMapper); - } - - @Override - public List> selectSpans(List traceIds) { - List gets = new ArrayList(traceIds.size()); - for (UUID traceId : traceIds) { - byte[] uuidBytes = BytesUtils.longLongToBytes(traceId.getMostSignificantBits(), traceId.getLeastSignificantBits()); - Get get = new Get(uuidBytes); - get.addFamily(COLFAM_SPAN); - gets.add(get); - } - return template2.get(HBaseTables.TRACES, gets, spanMapper); - } - - @Override - public List> selectSpans(Set traceIds) { - List gets = new ArrayList(traceIds.size()); - for (TraceId traceId : traceIds) { - Get get = new Get(traceId.getBytes()); - get.addFamily(COLFAM_SPAN); - gets.add(get); - } - return template2.get(HBaseTables.TRACES, gets, spanMapper); - } - - public List> selectSpansAndAnnotation(Set traceIds) { - List gets = new ArrayList(traceIds.size()); - for (TraceId traceId : traceIds) { - Get get = new Get(traceId.getBytes()); - get.addFamily(COLFAM_SPAN); - get.addFamily(COLFAM_ANNOTATION); - gets.add(get); - } - return template2.get(HBaseTables.TRACES, gets, spanAnnotationMapper); - } -} diff --git a/src/main/java/com/nhn/hippo/web/dao/HbaseTraceIndexDao.java b/src/main/java/com/nhn/hippo/web/dao/HbaseTraceIndexDao.java deleted file mode 100644 index 1a76ea8fe..000000000 --- a/src/main/java/com/nhn/hippo/web/dao/HbaseTraceIndexDao.java +++ /dev/null @@ -1,83 +0,0 @@ -package com.nhn.hippo.web.dao; - -import com.profiler.common.hbase.HBaseTables; -import com.profiler.common.hbase.HbaseOperations2; -import com.profiler.common.util.SpanUtils; -import org.apache.hadoop.hbase.client.Scan; -import org.apache.hadoop.hbase.util.Bytes; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.beans.factory.annotation.Qualifier; -import org.springframework.data.hadoop.hbase.RowMapper; -import org.springframework.stereotype.Repository; - -import java.util.ArrayList; -import java.util.List; - -/** - * - */ -@Repository -public class HbaseTraceIndexDao implements TraceIndexDao { - - private final Logger logger = LoggerFactory.getLogger(this.getClass()); - - private final byte[] COLFAM_TRACE = Bytes.toBytes("Trace"); - private final byte[] COLNAME_ID = Bytes.toBytes("ID"); - - @Autowired - private HbaseOperations2 hbaseOperations2; - - @Autowired - @Qualifier("traceIndexMapper") - private RowMapper traceIndexMapper; - - - private int scanCacheSize = 40; - - public void setScanCacheSize(int scanCacheSize) { - this.scanCacheSize = scanCacheSize; - } - - - @Override - public List scanTraceIndex(String agent, long start, long end) { - Scan scan = createScan(agent, start, end); - return hbaseOperations2.find(HBaseTables.TRACE_INDEX, scan, traceIndexMapper); - } - - @Override - public List> multiScanTraceIndex(String[] agents, long start, long end) { - final List multiScan = new ArrayList(agents.length); - for (String agent : agents) { - Scan scan = createScan(agent, start, end); - multiScan.add(scan); - } - return hbaseOperations2.find(HBaseTables.TRACE_INDEX, multiScan, traceIndexMapper); - } - - - private Scan createScan(String agent, long start, long end) { - - Scan scan = new Scan(); - // cache size를 지정해야 되는거 같음.?? - scan.setCaching(this.scanCacheSize); - - byte[] bAgent = Bytes.toBytes(agent); - byte[] traceIndexStartKey = SpanUtils.getTraceIndexRowKey(bAgent, start); - scan.setStartRow(traceIndexStartKey); -// TODO 추가 filter를 구현하여 scan시 중복된 값을 제가 할수 있음. 단 server에도 Filter 클래스가 배포되어야 한다. -// scan.setFilter(new ValueFilter()); - - byte[] traceIndexEndKey = SpanUtils.getTraceIndexRowKey(bAgent, end); - scan.setStopRow(traceIndexEndKey); - scan.addColumn(COLFAM_TRACE, COLNAME_ID); - scan.setId("traceIndexScan"); - - // json으로 변화해서 로그를 찍어서. 최초 변환 속도가 느림. - logger.debug("create scan:{}", scan); - return scan; - } - -} diff --git a/src/main/java/com/nhn/hippo/web/dao/hbase/HbaseApplicationIndexDao.java b/src/main/java/com/nhn/hippo/web/dao/hbase/HbaseApplicationIndexDao.java new file mode 100644 index 000000000..87547686e --- /dev/null +++ b/src/main/java/com/nhn/hippo/web/dao/hbase/HbaseApplicationIndexDao.java @@ -0,0 +1,54 @@ +package com.nhn.hippo.web.dao.hbase; + +import java.util.List; + +import org.apache.hadoop.hbase.client.Get; +import org.apache.hadoop.hbase.client.Scan; +import org.apache.hadoop.hbase.util.Bytes; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.data.hadoop.hbase.RowMapper; +import org.springframework.stereotype.Repository; + +import com.nhn.hippo.web.dao.ApplicationIndexDao; +import com.profiler.common.hbase.HBaseTables; +import com.profiler.common.hbase.HbaseOperations2; + +/** + * @author netspider + */ +@Repository +public class HbaseApplicationIndexDao implements ApplicationIndexDao { + + private final Logger logger = LoggerFactory.getLogger(this.getClass().getName()); + + @Autowired + private HbaseOperations2 hbaseOperations2; + + @Autowired + @Qualifier("applicationNameMapper") + private RowMapper applicationNameMapper; + + @Autowired + @Qualifier("agentIdMapper") + private RowMapper agentIdMapper; + + @Override + public List selectAllApplicationNames() { + Scan scan = new Scan(); + scan.setCaching(30); + return hbaseOperations2.find(HBaseTables.APPLICATION_INDEX, scan, applicationNameMapper); + } + + @Override + public String[] selectAgentIds(String applicationName) { + byte[] rowKey = Bytes.toBytes(applicationName); + + Get get = new Get(rowKey); + get.addFamily(HBaseTables.APPLICATION_INDEX_CF_AGENTS); + + return hbaseOperations2.get(HBaseTables.APPLICATION_INDEX, get, agentIdMapper); + } +} diff --git a/src/main/java/com/nhn/hippo/web/dao/hbase/HbaseApplicationTraceIndexDao.java b/src/main/java/com/nhn/hippo/web/dao/hbase/HbaseApplicationTraceIndexDao.java new file mode 100644 index 000000000..9eaa5a3e1 --- /dev/null +++ b/src/main/java/com/nhn/hippo/web/dao/hbase/HbaseApplicationTraceIndexDao.java @@ -0,0 +1,81 @@ +package com.nhn.hippo.web.dao.hbase; + +import java.util.ArrayList; +import java.util.List; + +import org.apache.hadoop.hbase.client.Scan; +import org.apache.hadoop.hbase.util.Bytes; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.data.hadoop.hbase.RowMapper; +import org.springframework.stereotype.Repository; + +import com.nhn.hippo.web.dao.ApplicationTraceIndexDao; +import com.profiler.common.hbase.HBaseTables; +import com.profiler.common.hbase.HbaseOperations2; +import com.profiler.common.util.SpanUtils; + +/** + * + */ +@Repository +public class HbaseApplicationTraceIndexDao implements ApplicationTraceIndexDao { + + private final Logger logger = LoggerFactory.getLogger(this.getClass()); + + private final byte[] COLFAM_TRACE = HBaseTables.APPLICATION_TRACE_INDEX_CF_TRACE; + private final byte[] COLNAME_ID = HBaseTables.APPLICATION_TRACE_INDEX_CN_ID; + + @Autowired + private HbaseOperations2 hbaseOperations2; + + @Autowired + @Qualifier("traceIndexMapper") + private RowMapper traceIndexMapper; + + private int scanCacheSize = 40; + + public void setScanCacheSize(int scanCacheSize) { + this.scanCacheSize = scanCacheSize; + } + + @Override + public List scanTraceIndex(String applicationName, long start, long end) { + Scan scan = createScan(applicationName, start, end); + return hbaseOperations2.find(HBaseTables.APPLICATION_TRACE_INDEX, scan, traceIndexMapper); + } + + @Override + public List> multiScanTraceIndex(String[] applicationNames, long start, long end) { + final List multiScan = new ArrayList(applicationNames.length); + for (String agent : applicationNames) { + Scan scan = createScan(agent, start, end); + multiScan.add(scan); + } + return hbaseOperations2.find(HBaseTables.APPLICATION_TRACE_INDEX, multiScan, traceIndexMapper); + } + + private Scan createScan(String agent, long start, long end) { + Scan scan = new Scan(); + // cache size를 지정해야 되는거 같음.?? + scan.setCaching(this.scanCacheSize); + + byte[] bAgent = Bytes.toBytes(agent); + byte[] traceIndexStartKey = SpanUtils.getTraceIndexRowKey(bAgent, start); + scan.setStartRow(traceIndexStartKey); + // TODO 추가 filter를 구현하여 scan시 중복된 값을 제가 할수 있음. 단 server에도 Filter 클래스가 + // 배포되어야 한다. + // scan.setFilter(new ValueFilter()); + + byte[] traceIndexEndKey = SpanUtils.getTraceIndexRowKey(bAgent, end); + scan.setStopRow(traceIndexEndKey); + scan.addColumn(COLFAM_TRACE, COLNAME_ID); + scan.setId("traceIndexScan"); + + // json으로 변화해서 로그를 찍어서. 최초 변환 속도가 느림. + logger.debug("create scan:{}", scan); + return scan; + } +} diff --git a/src/main/java/com/nhn/hippo/web/dao/hbase/HbaseRootTraceIndexDao.java b/src/main/java/com/nhn/hippo/web/dao/hbase/HbaseRootTraceIndexDao.java new file mode 100644 index 000000000..4350d2b55 --- /dev/null +++ b/src/main/java/com/nhn/hippo/web/dao/hbase/HbaseRootTraceIndexDao.java @@ -0,0 +1,85 @@ +package com.nhn.hippo.web.dao.hbase; + +import java.util.ArrayList; +import java.util.List; + +import org.apache.hadoop.hbase.client.Scan; +import org.apache.hadoop.hbase.util.Bytes; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.data.hadoop.hbase.RowMapper; +import org.springframework.stereotype.Repository; + +import com.nhn.hippo.web.dao.RootTraceIndexDao; +import com.profiler.common.hbase.HBaseTables; +import com.profiler.common.hbase.HbaseOperations2; +import com.profiler.common.util.SpanUtils; + +/** + * + */ +@Repository +public class HbaseRootTraceIndexDao implements RootTraceIndexDao { + + private Logger logger = LoggerFactory.getLogger(this.getClass()); + + private final byte[] COLFAM_TRACE = Bytes.toBytes("Trace"); + private final byte[] COLNAME_ID = Bytes.toBytes("ID"); + + @Autowired + private HbaseOperations2 hbaseOperations2; + + @Autowired + @Qualifier("traceIndexMapper") + private RowMapper traceIndexMapper; + + private int scanCacheSize = 40; + + public void setScanCacheSize(int scanCacheSize) { + this.scanCacheSize = scanCacheSize; + } + + @Override + public List scanTraceIndex(String agent, long start, long end) { + Scan scan = createScan(agent, start, end); + return hbaseOperations2.find(HBaseTables.ROOT_TRACE_INDEX, scan, traceIndexMapper); + } + + @Override + public List> multiScanTraceIndex(String[] agents, long start, long end) { + final List multiScan = new ArrayList(agents.length); + for (String agent : agents) { + Scan scan = createScan(agent, start, end); + multiScan.add(scan); + } + return hbaseOperations2.find(HBaseTables.ROOT_TRACE_INDEX, multiScan, traceIndexMapper); + } + + private Scan createScan(String agent, long start, long end) { + Scan scan = new Scan(); + // cache size를 지정해야 되는거 같음.?? + scan.setCaching(this.scanCacheSize); + + byte[] bAgent = Bytes.toBytes(agent); + byte[] traceIndexStartKey = SpanUtils.getTraceIndexRowKey(bAgent, start); + scan.setStartRow(traceIndexStartKey); + + byte[] traceIndexEndKey = SpanUtils.getTraceIndexRowKey(bAgent, end); + scan.setStopRow(traceIndexEndKey); + + scan.addColumn(COLFAM_TRACE, COLNAME_ID); + scan.setId("rootTraceIndexScan"); + + // json으로 변화해서 로그를 찍어서. 최초 변환 속도가 느림. + logger.debug("create scan:{}", scan); + return scan; + } + + @Override + public List parallelScanTraceIndex(String[] agents, long start, long end) { + return null; // To change body of implemented methods use File | + // Settings | File Templates. + } +} diff --git a/src/main/java/com/nhn/hippo/web/dao/hbase/HbaseTraceDao.java b/src/main/java/com/nhn/hippo/web/dao/hbase/HbaseTraceDao.java new file mode 100644 index 000000000..0d271cef9 --- /dev/null +++ b/src/main/java/com/nhn/hippo/web/dao/hbase/HbaseTraceDao.java @@ -0,0 +1,99 @@ +package com.nhn.hippo.web.dao.hbase; + +import java.util.ArrayList; +import java.util.List; +import java.util.Set; +import java.util.UUID; + +import org.apache.hadoop.hbase.client.Get; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.data.hadoop.hbase.RowMapper; +import org.springframework.stereotype.Repository; + +import com.nhn.hippo.web.dao.TraceDao; +import com.nhn.hippo.web.vo.TraceId; +import com.profiler.common.bo.SpanBo; +import com.profiler.common.hbase.HBaseTables; +import com.profiler.common.hbase.HbaseOperations2; +import com.profiler.common.util.BytesUtils; + +/** + * + */ +@Repository +public class HbaseTraceDao implements TraceDao { + + private final byte[] COLFAM_SPAN = HBaseTables.TRACES_CF_SPAN; + + private final byte[] COLFAM_ANNOTATION = HBaseTables.TRACES_CF_ANNOTATION; + + private Logger logger = LoggerFactory.getLogger(this.getClass()); + + @Autowired + private HbaseOperations2 template2; + + @Autowired + @Qualifier("spanMapper") + private RowMapper> spanMapper; + + @Autowired + @Qualifier("spanAnnotationMapper") + private RowMapper> spanAnnotationMapper; + + @Override + public List selectSpan(UUID traceId) { + byte[] uuidBytes = BytesUtils.longLongToBytes(traceId.getMostSignificantBits(), traceId.getLeastSignificantBits()); + return template2.get(HBaseTables.TRACES, uuidBytes, COLFAM_SPAN, spanMapper); + } + + public List selectSpanAndAnnotation(UUID traceId) { + byte[] uuidBytes = BytesUtils.longLongToBytes(traceId.getMostSignificantBits(), traceId.getLeastSignificantBits()); + Get get = new Get(uuidBytes); + get.addFamily(COLFAM_SPAN); + get.addFamily(COLFAM_ANNOTATION); + return template2.get(HBaseTables.TRACES, get, spanAnnotationMapper); + } + + @Override + public List selectSpan(long traceIdMost, long traceIdLeast) { + byte[] uuidBytes = BytesUtils.longLongToBytes(traceIdMost, traceIdLeast); + return template2.get(HBaseTables.TRACES, uuidBytes, COLFAM_SPAN, spanMapper); + } + + @Override + public List> selectSpans(List traceIds) { + List gets = new ArrayList(traceIds.size()); + for (UUID traceId : traceIds) { + byte[] uuidBytes = BytesUtils.longLongToBytes(traceId.getMostSignificantBits(), traceId.getLeastSignificantBits()); + Get get = new Get(uuidBytes); + get.addFamily(COLFAM_SPAN); + gets.add(get); + } + return template2.get(HBaseTables.TRACES, gets, spanMapper); + } + + @Override + public List> selectSpans(Set traceIds) { + List gets = new ArrayList(traceIds.size()); + for (TraceId traceId : traceIds) { + Get get = new Get(traceId.getBytes()); + get.addFamily(COLFAM_SPAN); + gets.add(get); + } + return template2.get(HBaseTables.TRACES, gets, spanMapper); + } + + public List> selectSpansAndAnnotation(Set traceIds) { + List gets = new ArrayList(traceIds.size()); + for (TraceId traceId : traceIds) { + Get get = new Get(traceId.getBytes()); + get.addFamily(COLFAM_SPAN); + get.addFamily(COLFAM_ANNOTATION); + gets.add(get); + } + return template2.get(HBaseTables.TRACES, gets, spanAnnotationMapper); + } +} diff --git a/src/main/java/com/nhn/hippo/web/dao/hbase/HbaseTraceIndexDao.java b/src/main/java/com/nhn/hippo/web/dao/hbase/HbaseTraceIndexDao.java new file mode 100644 index 000000000..17f77e3db --- /dev/null +++ b/src/main/java/com/nhn/hippo/web/dao/hbase/HbaseTraceIndexDao.java @@ -0,0 +1,81 @@ +package com.nhn.hippo.web.dao.hbase; + +import com.nhn.hippo.web.dao.TraceIndexDao; +import com.profiler.common.hbase.HBaseTables; +import com.profiler.common.hbase.HbaseOperations2; +import com.profiler.common.util.SpanUtils; +import org.apache.hadoop.hbase.client.Scan; +import org.apache.hadoop.hbase.util.Bytes; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.data.hadoop.hbase.RowMapper; +import org.springframework.stereotype.Repository; + +import java.util.ArrayList; +import java.util.List; + +/** + * + */ +@Repository +public class HbaseTraceIndexDao implements TraceIndexDao { + + private final Logger logger = LoggerFactory.getLogger(this.getClass()); + + private final byte[] COLFAM_TRACE = Bytes.toBytes("Trace"); + private final byte[] COLNAME_ID = Bytes.toBytes("ID"); + + @Autowired + private HbaseOperations2 hbaseOperations2; + + @Autowired + @Qualifier("traceIndexMapper") + private RowMapper traceIndexMapper; + + private int scanCacheSize = 40; + + public void setScanCacheSize(int scanCacheSize) { + this.scanCacheSize = scanCacheSize; + } + + @Override + public List scanTraceIndex(String agent, long start, long end) { + Scan scan = createScan(agent, start, end); + return hbaseOperations2.find(HBaseTables.TRACE_INDEX, scan, traceIndexMapper); + } + + @Override + public List> multiScanTraceIndex(String[] agents, long start, long end) { + final List multiScan = new ArrayList(agents.length); + for (String agent : agents) { + Scan scan = createScan(agent, start, end); + multiScan.add(scan); + } + return hbaseOperations2.find(HBaseTables.TRACE_INDEX, multiScan, traceIndexMapper); + } + + private Scan createScan(String agent, long start, long end) { + + Scan scan = new Scan(); + // cache size를 지정해야 되는거 같음.?? + scan.setCaching(this.scanCacheSize); + + byte[] bAgent = Bytes.toBytes(agent); + byte[] traceIndexStartKey = SpanUtils.getTraceIndexRowKey(bAgent, start); + scan.setStartRow(traceIndexStartKey); + // TODO 추가 filter를 구현하여 scan시 중복된 값을 제가 할수 있음. 단 server에도 Filter 클래스가 + // 배포되어야 한다. + // scan.setFilter(new ValueFilter()); + + byte[] traceIndexEndKey = SpanUtils.getTraceIndexRowKey(bAgent, end); + scan.setStopRow(traceIndexEndKey); + scan.addColumn(COLFAM_TRACE, COLNAME_ID); + scan.setId("traceIndexScan"); + + // json으로 변화해서 로그를 찍어서. 최초 변환 속도가 느림. + logger.debug("create scan:{}", scan); + return scan; + } +} diff --git a/src/main/java/com/nhn/hippo/web/service/FlowChartService.java b/src/main/java/com/nhn/hippo/web/service/FlowChartService.java index 3cd2720dc..6924f167a 100755 --- a/src/main/java/com/nhn/hippo/web/service/FlowChartService.java +++ b/src/main/java/com/nhn/hippo/web/service/FlowChartService.java @@ -30,6 +30,16 @@ public interface FlowChartService { */ public Set selectTraceIdsFromTraceIndex(String[] agentIds, long from, long to); + /** + * select traceIds from ApplicationTraceIndex table + * + * @param agentIds + * @param from + * @param to + * @return + */ + public Set selectTraceIdsFromApplicationTraceIndex(String applicationName, long from, long to); + /** * select call tree * @@ -52,4 +62,8 @@ public interface FlowChartService { * @return all of application names */ public List selectAllApplicationNames(); + + + public String[] selectAgentIds(String[] hosts); + } diff --git a/src/main/java/com/nhn/hippo/web/service/FlowChartServiceImpl.java b/src/main/java/com/nhn/hippo/web/service/FlowChartServiceImpl.java index fd604373b..6e98feb12 100755 --- a/src/main/java/com/nhn/hippo/web/service/FlowChartServiceImpl.java +++ b/src/main/java/com/nhn/hippo/web/service/FlowChartServiceImpl.java @@ -1,7 +1,10 @@ package com.nhn.hippo.web.service; +import java.util.ArrayList; import java.util.HashSet; +import java.util.Iterator; import java.util.List; +import java.util.Map; import java.util.Set; import org.slf4j.Logger; @@ -13,12 +16,16 @@ import org.springframework.stereotype.Service; import com.nhn.hippo.web.calltree.rpc.RPCCallTree; import com.nhn.hippo.web.calltree.server.ServerCallTree; import com.nhn.hippo.web.dao.ApplicationIndexDao; +import com.nhn.hippo.web.dao.ApplicationTraceIndexDao; import com.nhn.hippo.web.dao.RootTraceIndexDao; import com.nhn.hippo.web.dao.TraceDao; import com.nhn.hippo.web.dao.TraceIndexDao; import com.nhn.hippo.web.vo.TraceId; import com.profiler.common.bo.SpanBo; import com.profiler.common.hbase.HBaseClient; +import com.profiler.common.hbase.HBaseQuery; +import com.profiler.common.hbase.HBaseQuery.HbaseColumn; +import com.profiler.common.hbase.HBaseTables; /** * @author netspider @@ -44,6 +51,9 @@ public class FlowChartServiceImpl implements FlowChartService { @Autowired private ApplicationIndexDao applicationIndexDao; + @Autowired + private ApplicationTraceIndexDao applicationTraceIndexDao; + @Override public List selectAllApplicationNames() { return applicationIndexDao.selectAllApplicationNames(); @@ -111,4 +121,44 @@ public class FlowChartServiceImpl implements FlowChartService { } return tree.build(); } + + @Override + public Set selectTraceIdsFromApplicationTraceIndex(String applicationName, long from, long to) { + if (applicationName == null) { + throw new NullPointerException("applicationName"); + } + + if (logger.isTraceEnabled()) { + logger.trace("scan {}, {}, {}", new Object[] { applicationName, from, to }); + } + + List bytes = this.applicationTraceIndexDao.scanTraceIndex(applicationName, from, to); + Set result = new HashSet(); + for (byte[] traceId : bytes) { + TraceId tid = new TraceId(traceId); + result.add(tid); + logger.trace("traceid:{}", tid); + } + return result; + } + + @Override + public String[] selectAgentIds(String[] hosts) { + List column = new ArrayList(); + column.add(new HbaseColumn("Agents", "AgentID")); + + HBaseQuery query = new HBaseQuery(HBaseTables.APPLICATION_INDEX, null, null, column); + Iterator> iterator = client.getHBaseData(query); + + if (logger.isDebugEnabled()) { + while (iterator.hasNext()) { + logger.debug("selectedAgentId={}", iterator.next()); + } + logger.debug("!!!==============WARNING==============!!!"); + logger.debug("!!! selectAgentIds IS NOT IMPLEMENTED !!!"); + logger.debug("!!!===================================!!!"); + } + + return hosts; + } } diff --git a/src/main/resources/hbase.properties b/src/main/resources/hbase.properties index 80ba52846..28773bbdc 100644 --- a/src/main/resources/hbase.properties +++ b/src/main/resources/hbase.properties @@ -1,4 +1,4 @@ -#hbase.client.host=localhost -hbase.client.host=10.25.131.38 +hbase.client.host=localhost +#hbase.client.host=10.25.131.38 hbase.client.port=2181 hbase.htable.threads.max=64 \ No newline at end of file diff --git a/src/main/webapp/WEB-INF/views/flowserver.jsp b/src/main/webapp/WEB-INF/views/flowserver.jsp index 51451c85d..cb2d6c393 100644 --- a/src/main/webapp/WEB-INF/views/flowserver.jsp +++ b/src/main/webapp/WEB-INF/views/flowserver.jsp @@ -5,10 +5,10 @@ "nodes" : [ - { "name" : "${node.endPoint}" } + { "name" : "${node.endPoint}", "applicationName" : "${node.endPoint}" } - { "name" : "${node}" } + { "name" : "${node}", "applicationName" : "${node.applicationName}" } , diff --git a/src/main/webapp/index.html b/src/main/webapp/index.html index 81cf56397..f3a89093f 100644 --- a/src/main/webapp/index.html +++ b/src/main/webapp/index.html @@ -99,8 +99,8 @@ - - + + ~ diff --git a/src/test/java/com/nhn/hippo/web/service/SpanServiceTest.java b/src/test/java/com/nhn/hippo/web/service/SpanServiceTest.java index 4174fc90e..6b5632b84 100644 --- a/src/test/java/com/nhn/hippo/web/service/SpanServiceTest.java +++ b/src/test/java/com/nhn/hippo/web/service/SpanServiceTest.java @@ -6,7 +6,7 @@ import com.profiler.common.dto.thrift.Span; import com.profiler.common.hbase.HBaseTables; import com.profiler.common.hbase.HbaseTemplate2; import com.profiler.common.util.SpanUtils; -import com.profiler.server.dao.TraceDao; +import com.profiler.server.dao.Traces; import org.apache.hadoop.hbase.client.Delete; import org.apache.thrift.TException; import org.junit.Before; @@ -30,7 +30,7 @@ public class SpanServiceTest { @Autowired - private TraceDao traceDao; + private Traces traceDao; @Autowired private SpanService spanService; @@ -107,7 +107,7 @@ public class SpanServiceTest { private void insert(Span span) throws TException { - traceDao.insert(span); + traceDao.insert("JUNITApplicationName", span); } AtomicInteger id = new AtomicInteger(0);