From aad4d73eb2621900a3806496101a5d8a92f0d00c Mon Sep 17 00:00:00 2001 From: Woonduk Kang Date: Thu, 16 Aug 2012 10:54:43 +0000 Subject: [PATCH] =?UTF-8?q?[=EA=B0=95=EC=9A=B4=EB=8D=95]=20[LUCYSUS-1744]?= =?UTF-8?q?=20=EB=A9=94=EC=8B=9C=EC=A7=80=20=EB=A9=80=ED=8B=B0=ED=94=8C?= =?UTF-8?q?=EB=9E=99=EC=8B=B1=20=EA=B8=B0=EB=8A=A5=20=EB=B3=80=EA=B2=BD.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit git-svn-id: http://svn.bds.nhncorp.com/pe/hippo-tomcat-profiler/trunk@502 84d0f5b1-2673-498c-a247-62c4ff18d310 --- pom.xml | 14 -- src/main/java/com/profiler/dto/Constants.java | 38 ---- .../java/com/profiler/sender/DataSender.java | 166 ++++++++++-------- .../profiler/util/DefaultTBaseLocator.java | 32 +++- .../util/HeaderTBaseDeserializer.java | 2 +- .../java/com/profiler/util/TBaseLocator.java | 3 +- .../com/profiler/socket/UdpSocketTest.java | 13 ++ thrift/Header.thrift | 9 - thrift/MulplexedDTO.thrift | 31 ---- 9 files changed, 135 insertions(+), 173 deletions(-) delete mode 100644 src/main/java/com/profiler/dto/Constants.java delete mode 100644 thrift/Header.thrift delete mode 100644 thrift/MulplexedDTO.thrift diff --git a/pom.xml b/pom.xml index dd079c941..9abbced8f 100644 --- a/pom.xml +++ b/pom.xml @@ -242,20 +242,6 @@ - - - - - - - - - - - - - - diff --git a/src/main/java/com/profiler/dto/Constants.java b/src/main/java/com/profiler/dto/Constants.java deleted file mode 100644 index eac5f69b5..000000000 --- a/src/main/java/com/profiler/dto/Constants.java +++ /dev/null @@ -1,38 +0,0 @@ -/** - * Autogenerated by Thrift Compiler (0.8.0) - * - * DO NOT EDIT UNLESS YOU ARE SURE THAT YOU KNOW WHAT YOU ARE DOING - * @generated - */ -package com.profiler.dto; - -import org.apache.thrift.scheme.IScheme; -import org.apache.thrift.scheme.SchemeFactory; -import org.apache.thrift.scheme.StandardScheme; - -import org.apache.thrift.scheme.TupleScheme; -import org.apache.thrift.protocol.TTupleProtocol; -import java.util.List; -import java.util.ArrayList; -import java.util.Map; -import java.util.HashMap; -import java.util.EnumMap; -import java.util.Set; -import java.util.HashSet; -import java.util.EnumSet; -import java.util.Collections; -import java.util.BitSet; -import java.nio.ByteBuffer; -import java.util.Arrays; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -public class Constants { - - public static final short TYPE_JVM_INFO_DATA = (short)10; - - public static final short TYPE_REQUEST_DATA = (short)20; - - public static final short TYPE_REQUEST = (short)30; - -} diff --git a/src/main/java/com/profiler/sender/DataSender.java b/src/main/java/com/profiler/sender/DataSender.java index 8663cc24b..c8d9eaf82 100644 --- a/src/main/java/com/profiler/sender/DataSender.java +++ b/src/main/java/com/profiler/sender/DataSender.java @@ -6,101 +6,123 @@ import java.net.DatagramSocket; import java.net.InetSocketAddress; import java.net.SocketException; import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; import java.util.logging.Level; import java.util.logging.Logger; +import com.profiler.dto.Header; +import com.profiler.util.DefaultTBaseLocator; +import com.profiler.util.HeaderTBaseSerializer; +import com.profiler.util.TBaseLocator; import org.apache.thrift.TBase; import org.apache.thrift.TException; -import org.apache.thrift.TSerializer; -import org.apache.thrift.protocol.TBinaryProtocol; import com.profiler.config.TomcatProfilerConfig; -import com.profiler.dto.JVMInfoThriftDTO; -import com.profiler.dto.RequestDataListThriftDTO; -import com.profiler.dto.RequestThriftDTO; + /** - * * @author netspider - * */ public class DataSender extends Thread { - private final Logger logger = Logger.getLogger(DataSender.class.getName()); + private final Logger logger = Logger.getLogger(DataSender.class.getName()); - private final LinkedBlockingQueue> addedQueue = new LinkedBlockingQueue>(4096); + private final LinkedBlockingQueue> addedQueue = new LinkedBlockingQueue>(4096); - private final InetSocketAddress requestDataAddr = new InetSocketAddress(TomcatProfilerConfig.SERVER_IP, TomcatProfilerConfig.REQUEST_DATA_LISTEN_PORT); - private final InetSocketAddress requestTransactionDataAddr = new InetSocketAddress(TomcatProfilerConfig.SERVER_IP, TomcatProfilerConfig.REQUEST_TRANSACTION_DATA_LISTEN_PORT); - private final InetSocketAddress jvmDataAddr = new InetSocketAddress(TomcatProfilerConfig.SERVER_IP, TomcatProfilerConfig.JVM_DATA_LISTEN_PORT); + private final InetSocketAddress serverAddress = new InetSocketAddress(TomcatProfilerConfig.SERVER_IP, TomcatProfilerConfig.DEFUALT_PORT); - private static class SingletonHolder { - public static final DataSender INSTANCE = new DataSender(); - } + private DatagramSocket udpSocket = null; + private TBaseLocator locator = new DefaultTBaseLocator(); + // 주의 single thread용임 + private HeaderTBaseSerializer serializer = new HeaderTBaseSerializer(); - public static DataSender getInstance() { - return SingletonHolder.INSTANCE; - } + private DataSender() { + udpSocket = createSocket(); + setName("HIPPO-DataSender"); + setDaemon(true); + start(); + } - private DataSender() { - setName("HIPPO-DataSender"); - setDaemon(true); - start(); - } + private DatagramSocket createSocket() { + try { + DatagramSocket datagramSocket = new DatagramSocket(); + datagramSocket.setSoTimeout(1000 * 5); + datagramSocket.connect(serverAddress); + return datagramSocket; + } catch (SocketException e) { + return null; + } + } - public boolean addDataToSend(TBase data) { - // TODO: addedQueue가 full일 때 IllegalStateException처리. - return addedQueue.add(data); - } + public boolean addDataToSend(TBase data) { + // TODO: addedQueue가 full일 때 IllegalStateException처리. + return addedQueue.add(data); + } - // TODO: send timeout추 - // TODO: addedqueue에서 bulk로 drain - // TODO: sender thread가 한 개로 충분한가. - public void run() { - while (true) { - DatagramSocket udpSocket = null; - try { - TBase dto = addedQueue.take(); - // TODO TSerializer대신에 HeaderTBaseSerializer로 하고 Header를 생성하여 같이 넘기면 되며 - // 받는 쪽에서. HeaderTBaseDeSerialize로 받으면 됨. - // 단 header의 type 정보를 수동으로 넣어야 되는지가 불편함이 있음. - TSerializer serializer = new TSerializer(new TBinaryProtocol.Factory()); - byte[] sendData = serializer.serialize(dto); - // TODO: 포트 하나로 통일 시켜야 함. 일단 임시로 이렇게.. - InetSocketAddress address = null; - if (dto instanceof RequestDataListThriftDTO) { - address = requestDataAddr; - } else if (dto instanceof RequestThriftDTO) { - address = requestTransactionDataAddr; - } else if (dto instanceof JVMInfoThriftDTO) { - address = jvmDataAddr; - } + // TODO: sender thread가 한 개로 충분한가. + public void run() { + while (true) { + try { + TBase dto = take(); + if (dto == null) { + continue; + } + send(dto); + } catch (Exception e) { + logger.log(Level.WARNING, "Unexpected Error", e); + } + } + } - if (address == null) { - throw new IllegalArgumentException("Can't resolve receiver address."); - } + private void send(TBase dto) { + byte[] sendData = serialize(dto); + if (sendData == null) { + return; + } + DatagramPacket packet = new DatagramPacket(sendData, sendData.length); + if (udpSocket == null) { + // socket생성에 문제가 있으면 재생성? + udpSocket = createSocket(); + } + if (udpSocket != null) { + try { + udpSocket.send(packet); + if (logger.isLoggable(Level.FINE)) { + logger.fine("Data sent. " + dto); + } + } catch (IOException e) { + logger.log(Level.WARNING, "packet send error " + dto, e); + } + } + } - DatagramPacket packet = new DatagramPacket(sendData, sendData.length, address); - udpSocket = new DatagramSocket(); - udpSocket.send(packet); + // TODO: addedqueue에서 bulk로 drain + private TBase take() { + try { + return addedQueue.poll(5, TimeUnit.SECONDS); + } catch (InterruptedException e) { + return null; + } + } - //TODO: for test 로그레벨 바꾸기. - logger.info(String.format("Data sent. %s", dto)); - } catch (InterruptedException e) { - e.printStackTrace(); - } catch (TException e) { - e.printStackTrace(); - } catch (SocketException e) { - e.printStackTrace(); - } catch (IOException e) { - e.printStackTrace(); - } finally { - if (udpSocket != null) { - udpSocket.close(); - } - } - } - } + private byte[] serialize(TBase dto) { + Header header = createHeader(dto); + try { + return serializer.serialize(header, dto); + } catch (TException e) { + if (logger.isLoggable(Level.INFO)) { + logger.log(Level.INFO, "Serialize fail:" + dto, e); + } + return null; + } + } + + private Header createHeader(TBase dto) { + short type = locator.typeLookup(dto); + Header header = new Header(); + header.setType(type); + return header; + } } diff --git a/src/main/java/com/profiler/util/DefaultTBaseLocator.java b/src/main/java/com/profiler/util/DefaultTBaseLocator.java index d30d6d832..a2ebca0af 100644 --- a/src/main/java/com/profiler/util/DefaultTBaseLocator.java +++ b/src/main/java/com/profiler/util/DefaultTBaseLocator.java @@ -4,17 +4,35 @@ import com.profiler.dto.*; import org.apache.thrift.TBase; public class DefaultTBaseLocator implements TBaseLocator { + + private static final short JVM_INFO_THRIFT_DTO = 10; + private static final short REQUEST_DATA_LIST_THRIFT_DTO = 20; + private static final short REQUEST_THRIFT_DTO = 30; + + @Override - public TBase lookup(Header header) { - short type = header.getType(); - switch (type) { - case 10: + public TBase tBaseLookup(short type) { + switch (type) { + case JVM_INFO_THRIFT_DTO: return new JVMInfoThriftDTO(); - case 20: + case REQUEST_DATA_LIST_THRIFT_DTO: return new RequestDataListThriftDTO(); - case 30: + case REQUEST_THRIFT_DTO: return new RequestThriftDTO(); } - throw new IllegalArgumentException("Unsupported type:" + type + " " + header); + throw new IllegalArgumentException("Unsupported type:" + type); } + + public short typeLookup(TBase tbase) { + if(tbase instanceof JVMInfoThriftDTO) { + return JVM_INFO_THRIFT_DTO; + } + if(tbase instanceof RequestDataListThriftDTO) { + return REQUEST_DATA_LIST_THRIFT_DTO; + } + if(tbase instanceof RequestThriftDTO) { + return REQUEST_THRIFT_DTO; + } + throw new UnsupportedOperationException("Unsupported Type"); + } } diff --git a/src/main/java/com/profiler/util/HeaderTBaseDeserializer.java b/src/main/java/com/profiler/util/HeaderTBaseDeserializer.java index d98fc4438..516e5956f 100644 --- a/src/main/java/com/profiler/util/HeaderTBaseDeserializer.java +++ b/src/main/java/com/profiler/util/HeaderTBaseDeserializer.java @@ -43,7 +43,7 @@ public class HeaderTBaseDeserializer { trans_.reset(bytes); Header header = readHeader(); validate(header); - TBase base = locator.lookup(header); + TBase base = locator.tBaseLookup(header.getType()); base.read(protocol_); return base; } finally { diff --git a/src/main/java/com/profiler/util/TBaseLocator.java b/src/main/java/com/profiler/util/TBaseLocator.java index cc44b0942..a02a2a6be 100644 --- a/src/main/java/com/profiler/util/TBaseLocator.java +++ b/src/main/java/com/profiler/util/TBaseLocator.java @@ -4,5 +4,6 @@ import com.profiler.dto.Header; import org.apache.thrift.TBase; public interface TBaseLocator { - TBase lookup(Header header); + TBase tBaseLookup(short type); + short typeLookup(TBase tbase); } diff --git a/src/test/java/com/profiler/socket/UdpSocketTest.java b/src/test/java/com/profiler/socket/UdpSocketTest.java index cbc6dd0ba..3d26f26f3 100644 --- a/src/test/java/com/profiler/socket/UdpSocketTest.java +++ b/src/test/java/com/profiler/socket/UdpSocketTest.java @@ -154,4 +154,17 @@ public class UdpSocketTest { Thread.sleep(1000*3); } + +// @Test + public void createUdpSocket() throws IOException { + DatagramSocket so = new DatagramSocket(); +// so.bind(new InetSocketAddress("localhost", 8081)); +// DatagramSocket receiver = new DatagramSocket(new InetSocketAddress("localhost", 8082)); +// receiver.bind(new InetSocketAddress("localhost", 8082)); + + so.connect(new InetSocketAddress("localhost", 8082)); + so.send(new DatagramPacket(new byte[10], 10)); + +// receiver.receive(newDatagramPacket(1000)); + } } diff --git a/thrift/Header.thrift b/thrift/Header.thrift deleted file mode 100644 index 0f1944f09..000000000 --- a/thrift/Header.thrift +++ /dev/null @@ -1,9 +0,0 @@ -namespace java com.profiler.dto - -const byte SIGNATURE = 0xef; - -struct Header { - 1: required byte signature = SIGNATURE, - 2: required byte version = 0x10, - 3: required i16 type -} diff --git a/thrift/MulplexedDTO.thrift b/thrift/MulplexedDTO.thrift deleted file mode 100644 index 1f7677e14..000000000 --- a/thrift/MulplexedDTO.thrift +++ /dev/null @@ -1,31 +0,0 @@ -include "Header.thrift" -include "JVMInfoThriftDTO.thrift" -include "RequestDataThriftDTO.thrift" -include "RequestThriftDTO.thrift" - -namespace java com.profiler.dto - -const i16 TYPE_JVM_INFO_DATA = 10; -struct JVMInfoData { - 1: Header.Header header, - 2: JVMInfoThriftDTO.JVMInfoThriftDTO jvmInfoThriftDTO -} - -const i16 TYPE_REQUEST_DATA = 20; - -struct RequestDataList { - 1: Header.Header header, - 2: RequestDataThriftDTO.RequestDataListThriftDTO requestDataListThriftDTO -} - -const i16 TYPE_REQUEST = 30; -struct Request { - 1: Header.Header header, - 2: RequestThriftDTO.RequestThriftDTO request -} - - - - - -