From c6303b1f4d2b7ed63bc678444b7fdcbd2b9ac9fe Mon Sep 17 00:00:00 2001 From: Chisu Yu Date: Wed, 6 Mar 2013 05:59:25 +0000 Subject: [PATCH] =?UTF-8?q?[=EC=9C=A0=EC=B9=98=EC=88=98]=20[NOBTS]=20log?= =?UTF-8?q?=20message=20=EC=B6=94=EA=B0=80.?= 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@1295 84d0f5b1-2673-498c-a247-62c4ff18d310 --- src/main/java/com/profiler/Agent.java | 2 +- .../com/profiler/sender/UdpDataSender.java | 354 +++++++++--------- 2 files changed, 188 insertions(+), 168 deletions(-) diff --git a/src/main/java/com/profiler/Agent.java b/src/main/java/com/profiler/Agent.java index f8b85e481..ba7803165 100644 --- a/src/main/java/com/profiler/Agent.java +++ b/src/main/java/com/profiler/Agent.java @@ -180,7 +180,7 @@ public class Agent { agentInfo.setIsAlive(true); agentInfo.setTimestamp(this.startTime); - logger.info("Send startup information to HIPPO server. agentInfo=" + agentInfo); + logger.info("Send startup information to HIPPO server via " + this.priorityDataSender.getClass().getSimpleName() + ". agentInfo=" + agentInfo); send3(agentInfo); } diff --git a/src/main/java/com/profiler/sender/UdpDataSender.java b/src/main/java/com/profiler/sender/UdpDataSender.java index 0fe67649e..ef218e5dd 100644 --- a/src/main/java/com/profiler/sender/UdpDataSender.java +++ b/src/main/java/com/profiler/sender/UdpDataSender.java @@ -1,14 +1,5 @@ package com.profiler.sender; -import com.profiler.common.dto.Header; -import com.profiler.common.io.DefaultTBaseLocator; -import com.profiler.common.io.HeaderTBaseSerializer; -import com.profiler.common.io.TBaseLocator; -import com.profiler.context.Thriftable; -import com.profiler.util.Assert; -import org.apache.thrift.TBase; -import org.apache.thrift.TException; - import java.io.IOException; import java.net.DatagramPacket; import java.net.DatagramSocket; @@ -16,201 +7,230 @@ import java.net.InetSocketAddress; import java.net.SocketException; import java.util.ArrayList; import java.util.List; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.logging.Level; import java.util.logging.Logger; +import org.apache.thrift.TBase; +import org.apache.thrift.TException; + +import com.profiler.common.dto.Header; +import com.profiler.common.io.DefaultTBaseLocator; +import com.profiler.common.io.HeaderTBaseSerializer; +import com.profiler.common.io.TBaseLocator; +import com.profiler.context.Thriftable; +import com.profiler.util.Assert; + /** * @author netspider */ public class UdpDataSender implements DataSender, Runnable { - private final Logger logger = Logger.getLogger(UdpDataSender.class.getName()); + private final Logger logger = Logger.getLogger(UdpDataSender.class.getName()); - private final LinkedBlockingQueue queue = new LinkedBlockingQueue(1024); + private final LinkedBlockingQueue queue = new LinkedBlockingQueue(1024); - private int maxDrainSize = 10; - // 주의 single thread용임 - private List drain = new ArrayList(maxDrainSize); + private int maxDrainSize = 10; + // 주의 single thread용임 + private List drain = new ArrayList(maxDrainSize); - private DatagramSocket udpSocket = null; - private Thread ioThread; + private DatagramSocket udpSocket = null; + private Thread ioThread; - private TBaseLocator locator = new DefaultTBaseLocator(); - // 주의 single thread용임 - private HeaderTBaseSerializer serializer = new HeaderTBaseSerializer(); + private TBaseLocator locator = new DefaultTBaseLocator(); + // 주의 single thread용임 + private HeaderTBaseSerializer serializer = new HeaderTBaseSerializer(); + private AtomicBoolean allowInput = new AtomicBoolean(); + private CountDownLatch shutdownLatch = new CountDownLatch(1); - private AtomicBoolean started = new AtomicBoolean(); - private Object stopLock = new Object(); + public UdpDataSender(String host, int port) { + Assert.notNull(host, "host must not be null"); - public UdpDataSender(String host, int port) { - Assert.notNull(host, "host must not be null"); + // Socket 생성에 에러가 발생하면 Agent start가 안되게 변경. + this.udpSocket = createSocket(host, port); - // Socket 생성에 에러가 발생하면 Agent start가 안되게 변경. - this.udpSocket = createSocket(host, port); + this.allowInput.set(true); - this.ioThread = createIoThread(); + this.ioThread = createIoThread(); - this.started.set(true); - - logger.info("UdpDataSender initialized. host=" + host + ", port=" + port); - } + logger.info("UdpDataSender initialized. host=" + host + ", port=" + port); + } - private Thread createIoThread() { - Thread thread = new Thread(this); - thread.setName("HIPPO-UdpDataSender-IoThread"); - thread.setDaemon(true); - thread.start(); - return thread; - } + private Thread createIoThread() { + Thread thread = new Thread(this); + thread.setName("HIPPO-UdpDataSender-IoThread"); + thread.setDaemon(true); + thread.start(); + return thread; + } - private DatagramSocket createSocket(String host, int port) { - try { - DatagramSocket datagramSocket = new DatagramSocket(); - datagramSocket.setSoTimeout(1000 * 5); + private DatagramSocket createSocket(String host, int port) { + try { + DatagramSocket datagramSocket = new DatagramSocket(); + datagramSocket.setSoTimeout(1000 * 5); - InetSocketAddress serverAddress = new InetSocketAddress(host, port); - datagramSocket.connect(serverAddress); - return datagramSocket; - } catch (SocketException e) { - throw new IllegalStateException("DataramSocket create fail. Cause" + e.getMessage(), e); - } - } + InetSocketAddress serverAddress = new InetSocketAddress(host, port); + datagramSocket.connect(serverAddress); + return datagramSocket; + } catch (SocketException e) { + throw new IllegalStateException("DataramSocket create fail. Cause" + e.getMessage(), e); + } + } - public boolean send(TBase data) { - return putQueue(data); - } + public boolean send(TBase data) { + return putQueue(data); + } - public boolean send(Thriftable thriftable) { - return putQueue(thriftable); - } + public boolean send(Thriftable thriftable) { + return putQueue(thriftable); + } - private boolean putQueue(Object data) { - if (data == null) { - logger.warning("putQueue(). data is null"); - return false; - } - if (!started.get()) { - return false; - } - boolean offer = queue.offer(data); - if (!offer) { - if (logger.isLoggable(Level.WARNING)) { - logger.warning("Drop data. queue is full. size:" + queue.size()); - } - } - return offer; - } + private boolean putQueue(Object data) { + if (data == null) { + logger.warning("putQueue(). data is null"); + return false; + } + if (!allowInput.get()) { + return false; + } + boolean offer = queue.offer(data); + if (!offer) { + if (logger.isLoggable(Level.WARNING)) { + logger.warning("Drop data. queue is full. size:" + queue.size()); + } + } + return offer; + } + @Override + public void stop() { + allowInput.set(false); + + if (!isEmpty()) { + logger.info("Wait 5 seconds. Flushing queued data." + queue.size()); + } + + try { + shutdownLatch.await(5, TimeUnit.SECONDS); + } catch (InterruptedException e) { + logger.info("UdpDataSender stopped incompletely."); + } + + logger.info("UdpDataSender stopped."); + } - @Override - public void stop() { - if (!started.get()) { - return; - } - started.set(false); - // io thread 안전 종료. queue 비우기. - // TODO 종료 처리가 안이쁨. 고쳐야 될듯. - } + public void run() { + logger.info(Thread.currentThread().getName() + "(" + Thread.currentThread().getId() + ") started."); + doSend(); + } - public void run() { - doSend(); - } + private void doSend() { + drain: while (true) { + try { + if (!allowInput.get() && isEmpty()) { + shutdownLatch.countDown(); + break; + } - private void doSend() { - drain: - while (true) { - try { - List dtoList = takeN(); - if (dtoList != null) { - sendPacketN(dtoList); - continue; - } + List dtoList = takeN(); + if (dtoList != null) { + sendPacketN(dtoList); + continue; + } - while(true) { - Object dto = takeOne(); - if (dto != null) { - sendPacket(dto); - continue drain; - } - } + while (true) { + if (!allowInput.get() && isEmpty()) { + shutdownLatch.countDown(); + break; + } - } catch (Throwable th) { - logger.log(Level.WARNING, "Unexpected Error Cause:" + th.getMessage(), th); - } - } - } + Object dto = takeOne(); + if (dto != null) { + sendPacket(dto); + continue drain; + } + } + } catch (Throwable th) { + logger.log(Level.WARNING, "Unexpected Error Cause:" + th.getMessage(), th); + } + } + } - private void sendPacketN(List dtoList) { - for (Object dto : dtoList) { - try { - sendPacket(dto); - } catch (Throwable th) { - logger.log(Level.WARNING, "Unexpected Error Cause:" + th.getMessage(), th); - } - } - } + private void sendPacketN(List dtoList) { + for (Object dto : dtoList) { + try { + sendPacket(dto); + } catch (Throwable th) { + logger.log(Level.WARNING, "Unexpected Error Cause:" + th.getMessage(), th); + } + } + } - private void sendPacket(Object dto) { - TBase tBase; - if (dto instanceof TBase) { - tBase = (TBase) dto; - } else if (dto instanceof Thriftable) { - tBase = ((Thriftable) dto).toThrift(); - } else { - logger.warning("sendPacket fail. invalid type:" + dto.getClass()); - return; - } - byte[] sendData = serialize(tBase); - if (sendData == null) { - logger.warning("sendData is null"); - return; - } - DatagramPacket packet = new DatagramPacket(sendData, sendData.length); - 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); - } - } + private void sendPacket(Object dto) { + TBase tBase; + if (dto instanceof TBase) { + tBase = (TBase) dto; + } else if (dto instanceof Thriftable) { + tBase = ((Thriftable) dto).toThrift(); + } else { + logger.warning("sendPacket fail. invalid type:" + dto.getClass()); + return; + } + byte[] sendData = serialize(tBase); + if (sendData == null) { + logger.warning("sendData is null"); + return; + } + DatagramPacket packet = new DatagramPacket(sendData, sendData.length); + 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); + } + } - private Object takeOne() { - try { - return queue.poll(5, TimeUnit.SECONDS); - } catch (InterruptedException e) { - return null; - } - } + private Object takeOne() { + try { + return queue.poll(2, TimeUnit.SECONDS); + } catch (InterruptedException e) { + return null; + } + } - private List takeN() { - drain.clear(); - int size = queue.drainTo(drain, 10); - if (size <= 0) { - return null; - } - return drain; - } + private List takeN() { + drain.clear(); + int size = queue.drainTo(drain, 10); + if (size <= 0) { + return null; + } + return drain; + } - private byte[] serialize(TBase dto) { - try { - Header header = headerLookup(dto); - return serializer.serialize(header, dto); - } catch (TException e) { - if (logger.isLoggable(Level.WARNING)) { - logger.log(Level.WARNING, "Serialize fail:" + dto + " Caused:" + e.getMessage(), e); - } - return null; - } - } + private boolean isEmpty() { + return queue.size() == 0; + } - private Header headerLookup(TBase dto) throws TException { - // header 객체 생성을 안하고 정적 lookup이 되도록 변경. - return locator.headerLookup(dto); - } + private byte[] serialize(TBase dto) { + try { + Header header = headerLookup(dto); + return serializer.serialize(header, dto); + } catch (TException e) { + if (logger.isLoggable(Level.WARNING)) { + logger.log(Level.WARNING, "Serialize fail:" + dto + " Caused:" + e.getMessage(), e); + } + return null; + } + } + + private Header headerLookup(TBase dto) throws TException { + // header 객체 생성을 안하고 정적 lookup이 되도록 변경. + return locator.headerLookup(dto); + } }