package com.nhn.hippo.web.service; import java.util.ArrayList; import java.util.HashMap; import java.util.HashSet; import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Map.Entry; import java.util.Set; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import org.springframework.util.StopWatch; import com.nhn.hippo.web.calltree.rpc.RPCCallTree; import com.nhn.hippo.web.calltree.server.NodeIdGenerator; 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.TerminalStatisticsDao; import com.nhn.hippo.web.dao.TraceDao; import com.nhn.hippo.web.dao.TraceIndexDao; import com.nhn.hippo.web.vo.BusinessTransactions; import com.nhn.hippo.web.vo.TerminalStatistics; import com.nhn.hippo.web.vo.TraceId; import com.nhn.hippo.web.vo.scatter.Dot; import com.profiler.common.ServiceType; import com.profiler.common.bo.SpanBo; import com.profiler.common.bo.SubSpanBo; /** * @author netspider */ @Service public class FlowChartServiceImpl implements FlowChartService { private Logger logger = LoggerFactory.getLogger(this.getClass()); @Autowired private TraceDao traceDao; @Autowired private TraceIndexDao traceIndexDao; @Autowired private ApplicationIndexDao applicationIndexDao; @Autowired private ApplicationTraceIndexDao applicationTraceIndexDao; @Autowired private TerminalStatisticsDao terminalStatisticsDao; @Override public List selectAllApplicationNames() { return applicationIndexDao.selectAllApplicationNames(); } @Override public String[] selectAgentIdsFromApplicationName(String applicationName) { return applicationIndexDao.selectAgentIds(applicationName); } @Override public Set selectTraceIdsFromTraceIndex(String[] agentIds, long from, long to) { if (agentIds == null) { throw new NullPointerException("agentIds"); } if (agentIds.length == 1) { // single scan if (logger.isTraceEnabled()) { logger.trace("scan {}, {}, {}", new Object[] { agentIds[0], from, to }); } List> bytes = this.traceIndexDao.scanTraceIndex(agentIds[0], from, to); Set result = new HashSet(); for (List list : bytes) { for (byte[] traceId : list) { TraceId tid = new TraceId(traceId); result.add(tid); logger.trace("traceid:{}", tid); } } return result; } else { // multi scan 가능한 동일 open htable 에서 액세스함. List>> multiScan = this.traceIndexDao.multiScanTraceIndex(agentIds, from, to); Set result = new HashSet(); for (List> list : multiScan) { for (List scan : list) { for (byte[] traceId : scan) { result.add(new TraceId(traceId)); } } } return result; } } @Override public RPCCallTree selectRPCCallTree(Set traceIds) { final RPCCallTree tree = new RPCCallTree(); List> traces = this.traceDao.selectSpans(traceIds); for (List transaction : traces) { for (SpanBo eachTransaction : transaction) { tree.addSpan(eachTransaction); } } return tree.build(); } @Deprecated @Override public ServerCallTree selectServerCallTree(Set traceIds) { final ServerCallTree tree = new ServerCallTree(NodeIdGenerator.BY_APPLICATION_NAME); List> traces = this.traceDao.selectSpans(traceIds); for (List transaction : traces) { // List processed = refine(transaction); markRecursiveCall(transaction); for (SpanBo eachTransaction : transaction) { tree.addSpan(eachTransaction); } } return tree.build(); } /** * makes call tree of transaction detail view */ @Override public ServerCallTree selectServerCallTree(TraceId traceId) { final ServerCallTree tree = new ServerCallTree(NodeIdGenerator.BY_SERVER_INSTANCE); List transaction = this.traceDao.selectSpans(traceId); Set endPoints = new HashSet(); // List processed = refine(transaction); markRecursiveCall(transaction); for (SpanBo eachTransaction : transaction) { tree.addSpan(eachTransaction); endPoints.add(eachTransaction.getEndPoint()); } for (SpanBo eachTransaction : transaction) { List subSpanList = eachTransaction.getSubSpanList(); if (subSpanList == null) continue; for (SubSpanBo subTransaction : subSpanList) { // 통계정보로 잡지 않을 데이터는 스킵한다. if (!subTransaction.getServiceType().isRecordStatistics()) { continue; } // remove subspan of the rpc client if (!endPoints.contains(subTransaction.getEndPoint())) { // this is unknown cloud tree.addSubSpan(subTransaction); } } } return tree.build(); } /** * makes call tree of main view */ @Override public ServerCallTree selectServerCallTree(Set traceIds, String applicationName, long from, long to) { final Map terminalQueryParams = new HashMap(); final ServerCallTree tree = new ServerCallTree(NodeIdGenerator.BY_APPLICATION_NAME); StopWatch watch = new StopWatch(); watch.start("scanNonTerminalSpans"); // fetch non-terminal spans List> traces = this.traceDao.selectSpans(traceIds); watch.stop(); int totalNonTerminalSpansCount = 0; Set endPoints = new HashSet(); // processing spans for (List transaction : traces) { totalNonTerminalSpansCount += transaction.size(); // List processed = refine(transaction); markRecursiveCall(transaction); for (SpanBo eachTransaction : transaction) { tree.addSpan(eachTransaction); // make query param terminalQueryParams.put(eachTransaction.getServiceName(), eachTransaction.getServiceType()); endPoints.add(eachTransaction.getEndPoint()); } } if (logger.isInfoEnabled()) { logger.info("Fetch non-terminal spans elapsed : {}ms, {} traces, {} spans", new Object[] { watch.getLastTaskTimeMillis(), traces.size(), totalNonTerminalSpansCount }); } watch.start("scanTerminalStatistics"); // fetch terminal info for (Entry param : terminalQueryParams.entrySet()) { ServiceType svcType = param.getValue(); if (!svcType.isRpcClient() && !svcType.isUnknown() && !svcType.isTerminal()) { long start = System.currentTimeMillis(); List> terminals = terminalStatisticsDao.selectTerminal(param.getKey(), from, to); logger.info(" Fetch terminals of {} : {}ms", param.getKey(), System.currentTimeMillis() - start); for (Map terminal : terminals) { for (Entry entry : terminal.entrySet()) { // TODO 임시방편 TerminalStatistics t = entry.getValue(); if (!endPoints.contains(t.getTo())) { if (ServiceType.findServiceType(t.getToServiceType()).isRpcClient()) { t.setToServiceType(ServiceType.UNKNOWN_CLOUD.getCode()); } tree.addTerminalStatistics(t); } } } } } watch.stop(); logger.info("Fetch terminal statistics elapsed : {}ms", watch.getLastTaskTimeMillis()); return tree.build(); } @Deprecated private SpanBo findChildSpan(final List list, final SpanBo parent) { for (int i = 0; i < list.size(); i++) { SpanBo child = list.get(i); if (child.getParentSpanId() == parent.getSpanId()) { return child; } } return null; } @Deprecated private List refine(final List list) { for (int i = 0; i < list.size(); i++) { SpanBo span = list.get(i); if (span.getServiceType().isRpcClient()) { SpanBo child = findChildSpan(list, span); if (child != null) { child.setParentSpanId(span.getParentSpanId()); child.getAnnotationBoList().addAll(span.getAnnotationBoList()); list.remove(i); i--; continue; } else { // using as a terminal node. span.setServiceName(span.getEndPoint()); span.setServiceType(ServiceType.UNKNOWN_CLOUD); } } } return list; } private void markRecursiveCall(final List list) { for (int i = 0; i < list.size(); i++) { SpanBo a = list.get(i); for (int j = 0; j < list.size(); j++) { if (i == j) continue; SpanBo b = list.get(j); if (a.getServiceName().equals(b.getServiceName()) && a.getSpanId() == b.getParentSpanId()) { a.increaseRecursiveCallCount(); } } } } @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 (List list : bytes) { for (byte[] traceId : list) { TraceId tid = new TraceId(traceId); result.add(tid); logger.trace("traceid:{}", tid); } } return result; } @Deprecated @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; } @Override public List selectScatterData(String applicationName, long from, long to) { List> scanTrace = applicationTraceIndexDao.scanTraceScatter(applicationName, from, to); List list = new ArrayList(); for (List l : scanTrace) { for (Dot dot : l) { list.add(dot); } } return list; } @Override public List selectScatterData(String applicationName, long from, long to, int limit) { return applicationTraceIndexDao.scanTraceScatter2(applicationName, from, to, limit); } @Override public BusinessTransactions selectBusinessTransactions(Set traceIds, String applicationName, long from, long to) { List> traces = this.traceDao.selectSpans(traceIds); BusinessTransactions businessTransactions = new BusinessTransactions(); for (List transaction : traces) { for (SpanBo eachTransaction : transaction) { businessTransactions.add(eachTransaction); } } return businessTransactions; } }