diff --git a/src/main/java/com/profiler/context/Trace.java b/src/main/java/com/profiler/context/Trace.java index c38eaf485..fa7554b9b 100644 --- a/src/main/java/com/profiler/context/Trace.java +++ b/src/main/java/com/profiler/context/Trace.java @@ -3,11 +3,17 @@ package com.profiler.context; import java.util.logging.Level; import java.util.logging.Logger; +import com.profiler.Agent; import com.profiler.common.AnnotationNames; import com.profiler.common.ServiceType; +import com.profiler.common.dto.thrift.SqlMetaData; +import com.profiler.common.util.ParsingResult; +import com.profiler.common.util.SqlParser; import com.profiler.interceptor.MethodDescriptor; +import com.profiler.metadata.SqlCacheTable; import com.profiler.sender.DataSender; import com.profiler.sender.LoggingDataSender; +import com.profiler.util.Assert; /** * @author netspider @@ -30,6 +36,10 @@ public final class Trace { private Storage storage; + // TODO 아래 관련 핸들링 로직은 traceContext로 빼야지 이쁠뜻하다. + private SqlCacheTable sqlCacheTable; + private SqlParser sqlParser; + public Trace() { // traceObject에서 spanid의 유효성을 히스토리를 관리한다면 같은 thread에서는 span랜덤생성아이디의 충돌을 방지할수 있기는 함. TraceID traceId = TraceID.newTraceId(); @@ -220,23 +230,23 @@ public final class Trace { annotate(annotation.getCode()); } - public void recordException(Object result) { - if (result instanceof Throwable) { - Throwable th = (Throwable) result; - recordAttribute(AnnotationNames.EXCEPTION, th.getMessage()); + public void recordException(Object result) { + if (result instanceof Throwable) { + Throwable th = (Throwable) result; + recordAttribute(AnnotationNames.EXCEPTION, th.getMessage()); - try { - StackFrame currentStackFrame = getCurrentStackFrame(); - if (currentStackFrame instanceof RootStackFrame) { - ((RootStackFrame) currentStackFrame).getSpan().setException(true); - } else { - ((SubStackFrame) currentStackFrame).getSubSpan().setException(true); - } - } catch (Exception e) { - logger.log(Level.SEVERE, e.getMessage(), e); - } - } - } + try { + StackFrame currentStackFrame = getCurrentStackFrame(); + if (currentStackFrame instanceof RootStackFrame) { + ((RootStackFrame) currentStackFrame).getSpan().setException(true); + } else { + ((SubStackFrame) currentStackFrame).getSubSpan().setException(true); + } + } catch (Exception e) { + logger.log(Level.SEVERE, e.getMessage(), e); + } + } + } public void recordApi(MethodDescriptor methodDescriptor) { if (methodDescriptor == null) { @@ -274,6 +284,45 @@ public final class Trace { recordAttribute(key, (Object) value); } + public void recordSqlInfo(String sql) { + ParsingResult parsingResult = parseSql(sql); + recordAttribute(AnnotationNames.SQL_ID, parsingResult.getSql().hashCode()); + String output = parsingResult.getOutput(); + if (output != null && output.length() != 0) { + recordAttribute(AnnotationNames.SQL_PARAM, output); + } + } + + + public ParsingResult parseSql(String sql) { + Assert.notNull(sql, "sql must not null"); + + // 해당 api의 구현을 그냥 tarceContext api에 만들어야 될듯 하다. + ParsingResult parsingResult = this.sqlParser.normalizedSql(sql); + String normalizedSql = parsingResult.getSql(); + // 파싱시 변경되지 않았다면 동일 객체를 리턴하므로 그냥 ==비교를 하면 됨 + boolean newValue = this.sqlCacheTable.put(normalizedSql); + if (newValue) { + if (logger.isLoggable(Level.FINE)) { + // TODO hit% 로그를 남겨야 문제 발생시 도움이 될듯 하다. + logger.fine("NewSQLParsingResult:" + parsingResult); + } + // newValue란 의미는 cache에 인입됬다는 의미이고 이는 신규 sql문일 가능성이 있다는 의미임. + // 그러므로 메타데이터를 서버로 전송해야 한다. + + // 프로파일 데이터를 보내는데 사용되는 queue가 아니고, + // 좀더 급한 메시지만 별도 처리할수 있는 상대적으로 더 한가한 queue와 datasender를 별도로 가지고 있는게 좋을듯 하다. + SqlMetaData sqlMetaData = new SqlMetaData(); + sqlMetaData.setAgentId(Agent.getInstance().getAgentId()); + sqlMetaData.setStartTime(Agent.getInstance().getStartTime()); + sqlMetaData.setHashCode(normalizedSql.hashCode()); + sqlMetaData.setSql(normalizedSql); + this.getStorage().getDataSender().send(sqlMetaData); + } + // hashId그냥 return String에서 까보면 됨. + return parsingResult; + } + public void recordAttribute(final String key, final Object value) { if (!tracingEnabled) return; @@ -362,4 +411,13 @@ public final class Trace { } + public void setSqlCacheTable(SqlCacheTable sqlCacheTable) { + this.sqlCacheTable = sqlCacheTable; + } + + public void setSqlParser(SqlParser sqlParser) { + this.sqlParser = sqlParser; + } + + } \ No newline at end of file diff --git a/src/main/java/com/profiler/context/TraceContext.java b/src/main/java/com/profiler/context/TraceContext.java index 9bc962ec3..c478d8f34 100644 --- a/src/main/java/com/profiler/context/TraceContext.java +++ b/src/main/java/com/profiler/context/TraceContext.java @@ -1,8 +1,8 @@ package com.profiler.context; -import com.profiler.sender.DataSender; -import com.profiler.sender.LoggingDataSender; +import com.profiler.common.util.SqlParser; +import com.profiler.metadata.SqlCacheTable; import com.profiler.util.Assert; import com.profiler.util.NamedThreadLocal; @@ -37,6 +37,9 @@ public class TraceContext { private StorageFactory storageFactory; + private SqlCacheTable sqlTable = new SqlCacheTable(1000); + private SqlParser sqlParser = new SqlParser(); + public TraceContext() { } @@ -54,6 +57,8 @@ public class TraceContext { // trace.setDataSender(this.dataSender); Storage storage = storageFactory.createStorage(); trace.setStorage(storage); + trace.setSqlCacheTable(this.sqlTable); + trace.setSqlParser(this.sqlParser); // // trace.setTransactionId(transactionId.getAndIncrement()); threadLocal.set(trace); diff --git a/src/main/java/com/profiler/modifier/db/interceptor/PreparedStatementExecuteQueryInterceptor.java b/src/main/java/com/profiler/modifier/db/interceptor/PreparedStatementExecuteQueryInterceptor.java index 6873a286e..5c28a86f8 100644 --- a/src/main/java/com/profiler/modifier/db/interceptor/PreparedStatementExecuteQueryInterceptor.java +++ b/src/main/java/com/profiler/modifier/db/interceptor/PreparedStatementExecuteQueryInterceptor.java @@ -52,11 +52,14 @@ public class PreparedStatementExecuteQueryInterceptor implements StaticAroundInt trace.recordRpcName(databaseInfo.getType(), databaseInfo.getDatabaseId(), databaseInfo.getUrl()); trace.recordEndPoint(databaseInfo.getUrl()); String sql = getSql.invoke(target); - trace.recordAttribute(AnnotationNames.PREPAREDSTATEMENT, sql); + + // 일단 중복처리 + trace.recordAttribute(AnnotationNames.SQL, sql); + trace.recordSqlInfo(sql); Map bindValue = getBindValue.invoke(target); String bindString = toBindVariable(bindValue); - trace.recordAttribute(AnnotationNames.BINDVALUE, bindString); + trace.recordAttribute(AnnotationNames.SQL_BINDVALUE, bindString); // trace.recordApi(descriptor, args); trace.recordApi(apiId); diff --git a/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteQueryInterceptor.java b/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteQueryInterceptor.java index f2bc390d6..2c63b089a 100644 --- a/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteQueryInterceptor.java +++ b/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteQueryInterceptor.java @@ -1,5 +1,7 @@ package com.profiler.modifier.db.interceptor; +import com.profiler.common.AnnotationNames; +import com.profiler.common.util.ParsingResult; import com.profiler.context.Annotation; import com.profiler.context.Trace; import com.profiler.context.TraceContext; @@ -49,9 +51,7 @@ public class StatementExecuteQueryInterceptor implements StaticAroundInterceptor DatabaseInfo databaseInfo = (DatabaseInfo) this.getUrl.invoke(target); trace.recordRpcName(databaseInfo.getType(), databaseInfo.getDatabaseId(), databaseInfo.getUrl()); trace.recordEndPoint(databaseInfo.getUrl()); -// if (args.length > 0) { -// trace.recordAttribute("Statement", args[0]); -// } + } catch (Exception e) { if (logger.isLoggable(Level.WARNING)) { @@ -74,8 +74,17 @@ public class StatementExecuteQueryInterceptor implements StaticAroundInterceptor if (trace == null) { return; } - trace.recordApi(descriptor, args); + + trace.recordApi(descriptor); trace.recordException(result); + if (args.length > 0) { + Object arg = args[0]; + if (arg instanceof String) { + trace.recordSqlInfo((String) arg); + // TODO 일단 중복 처리. + trace.recordAttribute(AnnotationNames.SQL, args[0]); + } + } trace.markAfterTime(); trace.traceBlockEnd(); diff --git a/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteUpdateInterceptor.java b/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteUpdateInterceptor.java index 22d6a5907..3c630c859 100644 --- a/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteUpdateInterceptor.java +++ b/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteUpdateInterceptor.java @@ -1,6 +1,6 @@ package com.profiler.modifier.db.interceptor; -import com.profiler.context.Annotation; +import com.profiler.common.AnnotationNames; import com.profiler.context.Trace; import com.profiler.context.TraceContext; import com.profiler.interceptor.ApiIdSupport; @@ -48,7 +48,14 @@ public class StatementExecuteUpdateInterceptor implements StaticAroundIntercepto trace.recordRpcName(databaseInfo.getType(), databaseInfo.getDatabaseId(), databaseInfo.getUrl()); trace.recordEndPoint(databaseInfo.getUrl()); trace.recordApi(apiId, args); - + if (args.length > 0) { + Object arg = args[0]; + if (arg instanceof String) { + trace.recordSqlInfo((String) arg); + // TODO 일단 중복 처리. + trace.recordAttribute(AnnotationNames.SQL, args[0]); + } + } } catch (Exception e) { if (logger.isLoggable(Level.WARNING)) { diff --git a/src/main/java/com/profiler/sender/UdpDataSender.java b/src/main/java/com/profiler/sender/UdpDataSender.java index b0432a03f..6e649c5bc 100644 --- a/src/main/java/com/profiler/sender/UdpDataSender.java +++ b/src/main/java/com/profiler/sender/UdpDataSender.java @@ -95,6 +95,7 @@ public class UdpDataSender implements DataSender, Runnable { } started = false; // io thread 안전 종료. queue 비우기. + // TODO 종료 처리가 안이쁨. 고쳐야 될듯. } // TODO: sender thread가 한 개로 충분한가.