mirror of
https://github.com/wahyd4/pinpoint.git
synced 2026-08-17 08:46:22 +10:00
[강운덕] [LUCYSUS-1744] 리팩토링
git-svn-id: http://svn.bds.nhncorp.com/pe/hippo-tomcat-profiler/trunk@843 84d0f5b1-2673-498c-a247-62c4ff18d310
This commit is contained in:
@@ -1,132 +1,14 @@
|
||||
package com.profiler.sender;
|
||||
|
||||
import com.profiler.common.dto.Header;
|
||||
import com.profiler.common.util.DefaultTBaseLocator;
|
||||
import com.profiler.common.util.HeaderTBaseSerializer;
|
||||
import com.profiler.common.util.TBaseLocator;
|
||||
import com.profiler.config.ProfilerConfig;
|
||||
import org.apache.thrift.TBase;
|
||||
import org.apache.thrift.TException;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.net.DatagramPacket;
|
||||
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;
|
||||
|
||||
/**
|
||||
* @author netspider
|
||||
*/
|
||||
public class DataSender extends Thread {
|
||||
|
||||
private final Logger logger = Logger.getLogger(DataSender.class.getName());
|
||||
|
||||
private final LinkedBlockingQueue<TBase<?, ?>> addedQueue = new LinkedBlockingQueue<TBase<?, ?>>(4096);
|
||||
|
||||
private final InetSocketAddress serverAddress = new InetSocketAddress(ProfilerConfig.SERVER_IP, ProfilerConfig.SERVER_UDP_PORT);
|
||||
|
||||
private DatagramSocket udpSocket = null;
|
||||
private TBaseLocator locator = new DefaultTBaseLocator();
|
||||
// 주의 single thread용임
|
||||
private HeaderTBaseSerializer serializer = new HeaderTBaseSerializer();
|
||||
|
||||
private static class SingletonHolder {
|
||||
public static final DataSender INSTANCE = new DataSender();
|
||||
}
|
||||
|
||||
public static DataSender getInstance() {
|
||||
return SingletonHolder.INSTANCE;
|
||||
}
|
||||
|
||||
private DataSender() {
|
||||
udpSocket = createSocket();
|
||||
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);
|
||||
}
|
||||
|
||||
// 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TODO: addedqueue에서 bulk로 drain
|
||||
private TBase<?, ?> take() {
|
||||
try {
|
||||
return addedQueue.poll(5, TimeUnit.SECONDS);
|
||||
} catch (InterruptedException e) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
}
|
||||
package com.profiler.sender;
|
||||
|
||||
import org.apache.thrift.TBase;
|
||||
|
||||
/**
|
||||
*
|
||||
*/
|
||||
public interface DataSender {
|
||||
|
||||
boolean send(TBase<?, ?> data);
|
||||
|
||||
void stop();
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user