mirror of
https://github.com/wahyd4/pinpoint.git
synced 2026-08-24 04:06:34 +10:00
[유치수] [NOBTS] log message 추가.
git-svn-id: http://svn.bds.nhncorp.com/pe/hippo-tomcat-profiler/trunk@1295 84d0f5b1-2673-498c-a247-62c4ff18d310
This commit is contained in:
@@ -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);
|
||||
}
|
||||
|
||||
|
||||
@@ -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<Object> queue = new LinkedBlockingQueue<Object>(1024);
|
||||
private final LinkedBlockingQueue<Object> queue = new LinkedBlockingQueue<Object>(1024);
|
||||
|
||||
private int maxDrainSize = 10;
|
||||
// 주의 single thread용임
|
||||
private List<Object> drain = new ArrayList<Object>(maxDrainSize);
|
||||
private int maxDrainSize = 10;
|
||||
// 주의 single thread용임
|
||||
private List<Object> drain = new ArrayList<Object>(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<Object> dtoList = takeN();
|
||||
if (dtoList != null) {
|
||||
sendPacketN(dtoList);
|
||||
continue;
|
||||
}
|
||||
List<Object> 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<Object> 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<Object> 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<Object> takeN() {
|
||||
drain.clear();
|
||||
int size = queue.drainTo(drain, 10);
|
||||
if (size <= 0) {
|
||||
return null;
|
||||
}
|
||||
return drain;
|
||||
}
|
||||
private List<Object> 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);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user