diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/BaseUDPReceiver.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/BaseUDPHandlerFactory.java similarity index 67% rename from collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/BaseUDPReceiver.java rename to collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/BaseUDPHandlerFactory.java index 8a6c6335f..002870960 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/BaseUDPReceiver.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/BaseUDPHandlerFactory.java @@ -16,13 +16,16 @@ package com.navercorp.pinpoint.collector.receiver.udp; -import com.codahale.metrics.Timer; import com.navercorp.pinpoint.collector.receiver.DispatchHandler; import com.navercorp.pinpoint.collector.util.PacketUtils; import com.navercorp.pinpoint.thrift.io.*; import org.apache.thrift.TBase; import org.apache.thrift.TException; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.util.Assert; import java.io.IOException; import java.net.*; @@ -31,42 +34,51 @@ import java.net.*; * @author emeroad * @author netspider */ -public class BaseUDPReceiver extends AbstractUDPReceiver { +public class BaseUDPHandlerFactory implements PacketHandlerFactory, InitializingBean { + + private final Logger logger = LoggerFactory.getLogger(this.getClass()); + private DeserializerFactory deserializerFactory = new ThreadLocalHeaderTBaseDeserializerFactory(new HeaderTBaseDeserializerFactory()); - public BaseUDPReceiver(String receiverName, DispatchHandler dispatchHandler, String bindAddress, int port, int receiverBufferSize, int workerThreadSize, int workerThreadQueueSize) { - super(receiverName, dispatchHandler, bindAddress, port, receiverBufferSize, workerThreadSize, workerThreadQueueSize); + private UDPReceiver receiver; + private final DispatchHandler dispatchHandler; + + public BaseUDPHandlerFactory(DispatchHandler dispatchHandler) { + if (dispatchHandler == null) { + throw new NullPointerException("dispatchHandler must not be null"); + } + this.dispatchHandler = dispatchHandler; } - + + public void setReceiver(UDPReceiver receiver) { + this.receiver = receiver; + } + @Override - Runnable getPacketDispatcher(AbstractUDPReceiver receiver, DatagramPacket packet) { - return new DispatchPacket(receiver, packet); + public PacketHandler createPacketHandler() { + return new DispatchPacket(); } - private class DispatchPacket implements Runnable { - private final AbstractUDPReceiver receiver; - private final DatagramPacket packet; + @Override + public void afterPropertiesSet() throws Exception { + Assert.notNull(this.receiver, "receiver must not be null"); + } - private DispatchPacket(AbstractUDPReceiver receiver, DatagramPacket packet) { - if (packet == null) { - throw new NullPointerException("packet must not be null"); - } - this.receiver = receiver; - this.packet = packet; + private class DispatchPacket implements PacketHandler { + + private DispatchPacket() { } @Override - public void run() { - Timer.Context time = receiver.getTimer().time(); - - final HeaderTBaseDeserializer deserializer = (HeaderTBaseDeserializer) deserializerFactory.createDeserializer(); + public void receive(T packet) { + final HeaderTBaseDeserializer deserializer = deserializerFactory.createDeserializer(); TBase tBase = null; try { tBase = deserializer.deserialize(packet.getData()); if (tBase instanceof L4Packet) { if (logger.isDebugEnabled()) { - L4Packet packet = (L4Packet) tBase; - logger.debug("udp l4 packet {}", packet.getHeader()); + L4Packet l4Packet = (L4Packet) tBase; + logger.debug("udp l4 packet {}", l4Packet.getHeader()); } return; } @@ -75,11 +87,11 @@ public class BaseUDPReceiver extends AbstractUDPReceiver { if (logger.isDebugEnabled()) { logger.debug("received udp network availability check packet."); } - responseOK(); + responseOK(packet); return; } // dispatch signifies business logic execution - receiver.getDispatchHandler().dispatchSendMessage(tBase); + dispatchHandler.dispatchSendMessage(tBase); } catch (TException e) { if (logger.isWarnEnabled()) { logger.warn("packet serialize error. SendSocketAddress:{} Cause:{}", packet.getSocketAddress(), e.getMessage(), e); @@ -95,14 +107,10 @@ public class BaseUDPReceiver extends AbstractUDPReceiver { if (logger.isDebugEnabled()) { logger.debug("packet dump hex:{}", PacketUtils.dumpDatagramPacket(packet)); } - } finally { - receiver.getDatagramPacketPool().returnObject(packet); - // what should we do when an exception is thrown? - time.stop(); } } - private void responseOK() { + private void responseOK(DatagramPacket packet) { try { byte[] okBytes = NetworkAvailabilityCheckPacket.DATA_OK; DatagramPacket pongPacket = new DatagramPacket(okBytes, okBytes.length, packet.getSocketAddress()); diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/ChunkedUDPReceiver.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/ChunkedUDPPacketHandlerFactory.java similarity index 71% rename from collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/ChunkedUDPReceiver.java rename to collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/ChunkedUDPPacketHandlerFactory.java index 904454ed7..dab5dfb7f 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/ChunkedUDPReceiver.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/ChunkedUDPPacketHandlerFactory.java @@ -16,13 +16,16 @@ package com.navercorp.pinpoint.collector.receiver.udp; -import com.codahale.metrics.Timer; import com.navercorp.pinpoint.collector.receiver.DispatchHandler; import com.navercorp.pinpoint.collector.util.PacketUtils; import com.navercorp.pinpoint.thrift.io.*; import org.apache.thrift.TBase; import org.apache.thrift.TException; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.util.Assert; import java.io.IOException; import java.net.*; @@ -33,36 +36,44 @@ import java.util.List; * * @author jaehong.kim */ -public class ChunkedUDPReceiver extends AbstractUDPReceiver { +public class ChunkedUDPPacketHandlerFactory implements PacketHandlerFactory, InitializingBean { + + private final Logger logger = LoggerFactory.getLogger(this.getClass()); private final DeserializerFactory deserializerFactory = new ThreadLocalHeaderTBaseDeserializerFactory(new ChunkHeaderTBaseDeserializerFactory()); - - public ChunkedUDPReceiver(String receiverName, DispatchHandler dispatchHandler, String bindAddress, int port, int receiverBufferSize, int workerThreadSize, int workerThreadQueueSize) { - super(receiverName, dispatchHandler, bindAddress, port, receiverBufferSize, workerThreadSize, workerThreadQueueSize); + private UDPReceiver receiver; + private final DispatchHandler dispatchHandler; + + public ChunkedUDPPacketHandlerFactory(DispatchHandler dispatchHandler) { + if (dispatchHandler == null) { + throw new NullPointerException("dispatchHandler must not be null"); + } + this.dispatchHandler = dispatchHandler; } - + + public void setReceiver(UDPReceiver receiver) { + this.receiver = receiver; + } + + @Override - Runnable getPacketDispatcher(AbstractUDPReceiver receiver, DatagramPacket packet) { - return new DispatchPacket(receiver, packet); + public void afterPropertiesSet() throws Exception { + Assert.notNull(this.receiver, "receiver must not be null"); } - private class DispatchPacket implements Runnable { - private final AbstractUDPReceiver receiver; - private final DatagramPacket packet; + @Override + public PacketHandler createPacketHandler() { + return new DispatchPacket(); + } - private DispatchPacket(AbstractUDPReceiver receiver, DatagramPacket packet) { - if (packet == null) { - throw new NullPointerException("packet must not be null"); - } - this.receiver = receiver; - this.packet = packet; + private class DispatchPacket implements PacketHandler { + + private DispatchPacket() { } @Override - public void run() { - Timer.Context time = receiver.getTimer().time(); - + public void receive(T packet) { final ChunkHeaderTBaseDeserializer deserializer = deserializerFactory.createDeserializer(); try { List> list = deserializer.deserialize(packet.getData(), packet.getOffset(), packet.getLength()); @@ -73,8 +84,8 @@ public class ChunkedUDPReceiver extends AbstractUDPReceiver { for (TBase tBase : list) { if (tBase instanceof L4Packet) { if (logger.isDebugEnabled()) { - L4Packet packet = (L4Packet) tBase; - logger.debug("udp l4 packet {}", packet.getHeader()); + L4Packet l4Packet = (L4Packet) tBase; + logger.debug("udp l4 packet {}", l4Packet.getHeader()); } continue; } @@ -83,11 +94,11 @@ public class ChunkedUDPReceiver extends AbstractUDPReceiver { if (logger.isDebugEnabled()) { logger.debug("received udp network availability check packet."); } - responseOK(); + responseOK(packet); continue; } // dispatch signifies business logic execution - receiver.getDispatchHandler().dispatchSendMessage(tBase); + dispatchHandler.dispatchSendMessage(tBase); } } catch (TException e) { if (logger.isWarnEnabled()) { @@ -103,13 +114,10 @@ public class ChunkedUDPReceiver extends AbstractUDPReceiver { if (logger.isDebugEnabled()) { logger.debug("packet dump hex:{}", PacketUtils.dumpDatagramPacket(packet)); } - } finally { - receiver.getDatagramPacketPool().returnObject(packet); - time.stop(); } } - private void responseOK() { + private void responseOK(DatagramPacket packet) { try { byte[] okBytes = NetworkAvailabilityCheckPacket.DATA_OK; DatagramPacket pongPacket = new DatagramPacket(okBytes, okBytes.length, packet.getSocketAddress()); diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/PacketHandler.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/PacketHandler.java new file mode 100644 index 000000000..88fbf1dee --- /dev/null +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/PacketHandler.java @@ -0,0 +1,24 @@ +/* + * Copyright 2014 NAVER Corp. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.navercorp.pinpoint.collector.receiver.udp; + +/** + * @author emeroad + */ +public interface PacketHandler { + void receive(T packet); +} diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/PacketHandlerFactory.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/PacketHandlerFactory.java new file mode 100644 index 000000000..527f1d2a2 --- /dev/null +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/PacketHandlerFactory.java @@ -0,0 +1,24 @@ +/* + * Copyright 2014 NAVER Corp. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.navercorp.pinpoint.collector.receiver.udp; + +/** + * @author emeroad + */ +public interface PacketHandlerFactory { + PacketHandler createPacketHandler(); +} diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/PooledPacketWrap.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/PooledPacketWrap.java new file mode 100644 index 000000000..2bcbc36e0 --- /dev/null +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/PooledPacketWrap.java @@ -0,0 +1,50 @@ +/* + * Copyright 2014 NAVER Corp. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.navercorp.pinpoint.collector.receiver.udp; + +import com.navercorp.pinpoint.collector.util.PooledObject; + +import java.net.DatagramPacket; + +/** + * @author emeroad + */ +public class PooledPacketWrap implements Runnable { + private final PacketHandler packetHandler; + private final PooledObject pooledObject; + + public PooledPacketWrap(PacketHandler packetHandler, PooledObject pooledObject) { + if (packetHandler == null) { + throw new NullPointerException("packetReceiveHandler must not be null"); + } + if (pooledObject == null) { + throw new NullPointerException("pooledObject must not be null"); + } + this.packetHandler = packetHandler; + this.pooledObject = pooledObject; + } + + @Override + public void run() { + final DatagramPacket packet = pooledObject.getObject(); + try { + packetHandler.receive(packet); + } catch (Exception e) { + pooledObject.returnObject(); + } + } +} diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/SpanStreamUDPReceiver.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/SpanStreamUDPPacketHandlerFactory.java similarity index 76% rename from collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/SpanStreamUDPReceiver.java rename to collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/SpanStreamUDPPacketHandlerFactory.java index c546a87c0..c35f9567c 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/SpanStreamUDPReceiver.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/SpanStreamUDPPacketHandlerFactory.java @@ -21,10 +21,9 @@ import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.List; +import com.navercorp.pinpoint.collector.receiver.DispatchHandler; import org.apache.thrift.TBase; -import com.codahale.metrics.Timer; -import com.navercorp.pinpoint.collector.receiver.DispatchHandler; import com.navercorp.pinpoint.thrift.dto.TSpan; import com.navercorp.pinpoint.thrift.dto.TSpanChunk; import com.navercorp.pinpoint.thrift.dto.TSpanEvent; @@ -33,41 +32,46 @@ import com.navercorp.pinpoint.thrift.io.HeaderTBaseDeserializer; import com.navercorp.pinpoint.thrift.io.HeaderTBaseDeserializerFactory; import com.navercorp.pinpoint.thrift.io.SpanStreamConstants; import com.navercorp.pinpoint.thrift.io.ThreadLocalHeaderTBaseDeserializerFactory; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.util.Assert; /** * @author Taejin Koo */ -public class SpanStreamUDPReceiver extends AbstractUDPReceiver { +public class SpanStreamUDPPacketHandlerFactory implements PacketHandlerFactory, InitializingBean { + + private final Logger logger = LoggerFactory.getLogger(this.getClass()); private DeserializerFactory deserializerFactory = new ThreadLocalHeaderTBaseDeserializerFactory(new HeaderTBaseDeserializerFactory()); + private final DispatchHandler dispatchHandler; - public SpanStreamUDPReceiver(String receiverName, DispatchHandler dispatchHandler, String bindAddress, int port, int receiverBufferSize, - int workerThreadSize, int workerThreadQueueSize) { - super(receiverName, dispatchHandler, bindAddress, port, receiverBufferSize, workerThreadSize, workerThreadQueueSize); + public SpanStreamUDPPacketHandlerFactory(DispatchHandler dispatchHandler) { + if (dispatchHandler == null) { + throw new NullPointerException("dispatchHandler must not be null"); + } + this.dispatchHandler = dispatchHandler; } @Override - Runnable getPacketDispatcher(AbstractUDPReceiver receiver, DatagramPacket packet) { - return new DispatchPacket(receiver, packet); + public void afterPropertiesSet() throws Exception { + Assert.notNull(this.dispatchHandler, "dispatchHandler must not be null"); } - private class DispatchPacket implements Runnable { - private final AbstractUDPReceiver receiver; - private final DatagramPacket packet; + @Override + public PacketHandler createPacketHandler() { + return new DispatchPacket(); + } - private DispatchPacket(AbstractUDPReceiver receiver, DatagramPacket packet) { - if (packet == null) { - throw new NullPointerException("packet must not be null"); - } - this.receiver = receiver; - this.packet = packet; + private class DispatchPacket implements PacketHandler { + + private DispatchPacket() { } @Override - public void run() { - Timer.Context time = receiver.getTimer().time(); - - final HeaderTBaseDeserializer deserializer = (HeaderTBaseDeserializer) deserializerFactory.createDeserializer(); + public void receive(DatagramPacket packet) { + final HeaderTBaseDeserializer deserializer = deserializerFactory.createDeserializer(); ByteBuffer requestBuffer = ByteBuffer.wrap(packet.getData()); if (requestBuffer.remaining() < SpanStreamConstants.START_PROTOCOL_BUFFER_SIZE) { @@ -106,14 +110,10 @@ public class SpanStreamUDPReceiver extends AbstractUDPReceiver { } else if (tBase instanceof TSpanChunk) { ((TSpanChunk) tBase).setSpanEventList(spanEventList); } - receiver.getDispatchHandler().dispatchRequestMessage(tBase); + dispatchHandler.dispatchRequestMessage(tBase); } } catch (Exception e) { logger.warn("Failed to handle receive packet.", e); - } finally { - receiver.getDatagramPacketPool().returnObject(packet); - // what should we do when an exception is thrown? - time.stop(); } } } diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/TimingWrap.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/TimingWrap.java new file mode 100644 index 000000000..d8f2035d5 --- /dev/null +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/TimingWrap.java @@ -0,0 +1,49 @@ +/* + * Copyright 2014 NAVER Corp. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.navercorp.pinpoint.collector.receiver.udp; + +import com.codahale.metrics.Timer; + +/** + * @author emeroad + */ +public class TimingWrap implements Runnable { + private final Timer timer; + private final Runnable child; + + public TimingWrap(Timer timer, Runnable child) { + if (timer == null) { + throw new NullPointerException("timer must not be null"); + } + if (child == null) { + throw new NullPointerException("child must not be null"); + } + this.timer = timer; + this.child = child; + } + + @Override + public void run() { + final Timer.Context time = timer.time(); + try { + child.run(); + } finally { + time.stop(); + } + + } +} diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/AbstractUDPReceiver.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/UDPReceiver.java similarity index 80% rename from collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/AbstractUDPReceiver.java rename to collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/UDPReceiver.java index 37abc3b5a..da5bc5e0b 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/AbstractUDPReceiver.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/udp/UDPReceiver.java @@ -20,10 +20,7 @@ import com.codahale.metrics.Counter; import com.codahale.metrics.MetricRegistry; import com.codahale.metrics.Timer; import com.navercorp.pinpoint.collector.receiver.DataReceiver; -import com.navercorp.pinpoint.collector.receiver.DispatchHandler; -import com.navercorp.pinpoint.collector.util.DatagramPacketFactory; -import com.navercorp.pinpoint.collector.util.ObjectPool; -import com.navercorp.pinpoint.collector.util.PacketUtils; +import com.navercorp.pinpoint.collector.util.*; import com.navercorp.pinpoint.common.util.ExecutorFactory; import com.navercorp.pinpoint.common.util.PinpointThreadFactory; import com.navercorp.pinpoint.rpc.util.CpuUtils; @@ -46,11 +43,10 @@ import java.util.concurrent.atomic.AtomicInteger; * @author emeroad * @author netspider * @author jaehong.kim - * change to abstract class */ -public abstract class AbstractUDPReceiver implements DataReceiver { +public class UDPReceiver implements DataReceiver { - protected final Logger logger = LoggerFactory.getLogger(this.getClass().getName()); + private final Logger logger; private String bindAddress; private int port; @@ -81,44 +77,50 @@ public abstract class AbstractUDPReceiver implements DataReceiver { private volatile DatagramSocket socket = null; - private DispatchHandler dispatchHandler; - + private PacketHandlerFactory packetHandlerFactory; private AtomicInteger rejectedExecutionCount = new AtomicInteger(0); private AtomicBoolean state = new AtomicBoolean(true); - public AbstractUDPReceiver() { + public UDPReceiver() { + this.logger = LoggerFactory.getLogger(this.getClass()); } - public AbstractUDPReceiver(String receiverName, DispatchHandler dispatchHandler, String bindAddress, int port, int receiverBufferSize, int workerThreadSize, int workerThreadQueueSize) { - if (dispatchHandler == null) { - throw new NullPointerException("dispatchHandler must not be null"); + public UDPReceiver(String receiverName, PacketHandlerFactory packetHandlerFactory, String bindAddress, int port, int receiverBufferSize, int workerThreadSize, int workerThreadQueueSize) { + if (receiverName != null) { + this.logger = LoggerFactory.getLogger(receiverName); + } else { + this.logger = LoggerFactory.getLogger(this.getClass()); + } + if (packetHandlerFactory == null) { + throw new NullPointerException("packetHandlerFactory must not be null"); } if (bindAddress == null) { throw new NullPointerException("bindAddress must not be null"); } + + this.receiverName = receiverName; - this.dispatchHandler = dispatchHandler; this.bindAddress = bindAddress; this.port = port; this.receiverBufferSize = receiverBufferSize; this.workerThreadSize = workerThreadSize; this.workerThreadQueueSize = workerThreadQueueSize; + this.packetHandlerFactory = packetHandlerFactory; } - abstract Runnable getPacketDispatcher(AbstractUDPReceiver receiver, DatagramPacket packet); - + public void afterPropertiesSet() { - Assert.notNull(dispatchHandler, "dispatchHandler must not be null"); Assert.notNull(metricRegistry, "metricRegistry must not be null"); + Assert.notNull(packetHandlerFactory, "packetHandlerFactory must not be null"); this.socket = createSocket(bindAddress, port, receiverBufferSize); final int packetPoolSize = getPacketPoolSize(workerThreadSize, workerThreadQueueSize); - this.datagramPacketPool = new ObjectPool(new DatagramPacketFactory(), packetPoolSize); + this.datagramPacketPool = new DefaultObjectPool(new DatagramPacketFactory(), packetPoolSize); this.worker = ExecutorFactory.newFixedThreadPool(workerThreadSize, workerThreadQueueSize, receiverName + "-Worker", true); this.timer = metricRegistry.timer(receiverName + "-timer"); @@ -136,10 +138,11 @@ public abstract class AbstractUDPReceiver implements DataReceiver { // need shutdown logic while (state.get()) { - DatagramPacket packet = read0(); - if (packet == null) { + PooledObject pooledPacket = read0(); + if (pooledPacket == null) { continue; } + DatagramPacket packet = pooledPacket.getObject(); if (packet.getLength() == 0) { if (debugEnabled) { logger.debug("length is 0 ip:{}, port:{}", packet.getAddress(), packet.getPort()); @@ -150,7 +153,8 @@ public abstract class AbstractUDPReceiver implements DataReceiver { logger.debug("pool getActiveCount:{}", worker.getActiveCount()); } try { - worker.execute(getPacketDispatcher(this, packet)); + Runnable dispatchTask = wrapDispatchTask(pooledPacket); + worker.execute(dispatchTask); } catch (RejectedExecutionException ree) { rejectedCounter.inc(); final int error = rejectedExecutionCount.incrementAndGet(); @@ -165,13 +169,20 @@ public abstract class AbstractUDPReceiver implements DataReceiver { } } - private DatagramPacket read0() { + private Runnable wrapDispatchTask(PooledObject pooledPacket) { + PacketHandler dispatchPacket = packetHandlerFactory.createPacketHandler(); + PooledPacketWrap pooledPacketWrap = new PooledPacketWrap(dispatchPacket, pooledPacket); + return new TimingWrap(this.timer, pooledPacketWrap); + } + + private PooledObject read0() { boolean success = false; - DatagramPacket packet = datagramPacketPool.getObject(); - if (packet == null) { + PooledObject pooledObject = datagramPacketPool.getObject(); + if (pooledObject == null) { logger.error("datagramPacketPool is empty"); return null; } + DatagramPacket packet = pooledObject.getObject(); try { try { socket.receive(packet); @@ -195,10 +206,10 @@ public abstract class AbstractUDPReceiver implements DataReceiver { return null; } finally { if (!success) { - datagramPacketPool.returnObject(packet); + pooledObject.returnObject(); } } - return packet; + return pooledObject; } private DatagramSocket createSocket(String bindAddress, int port, int receiveBufferSize) { @@ -266,18 +277,6 @@ public abstract class AbstractUDPReceiver implements DataReceiver { } } - public Timer getTimer() { - return timer; - } - - public DispatchHandler getDispatchHandler() { - return dispatchHandler; - } - - public ObjectPool getDatagramPacketPool() { - return datagramPacketPool; - } - public DatagramSocket getSocket() { return socket; } diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/util/ByteBufferFactory.java b/collector/src/main/java/com/navercorp/pinpoint/collector/util/ByteBufferFactory.java new file mode 100644 index 000000000..430301e48 --- /dev/null +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/util/ByteBufferFactory.java @@ -0,0 +1,37 @@ +/* + * Copyright 2014 NAVER Corp. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.navercorp.pinpoint.collector.util; + +import java.nio.ByteBuffer; + +/** + * @author emeroad + */ +public class ByteBufferFactory implements ObjectPoolFactory { + + private static final int AcceptedSize = 65507; + + @Override + public ByteBuffer create() { + return ByteBuffer.allocate(AcceptedSize); + } + + @Override + public void beforeReturn(ByteBuffer packet) { + packet.clear(); + } +} diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/util/DefaultObjectPool.java b/collector/src/main/java/com/navercorp/pinpoint/collector/util/DefaultObjectPool.java new file mode 100644 index 000000000..3298478f6 --- /dev/null +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/util/DefaultObjectPool.java @@ -0,0 +1,93 @@ +/* + * Copyright 2014 NAVER Corp. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.navercorp.pinpoint.collector.util; + +import java.util.Queue; +import java.util.concurrent.ConcurrentLinkedQueue; + +/** + * @author emeroad + */ +public class DefaultObjectPool implements ObjectPool { + + // you don't need a blocking queue. There must be enough objects in a queue. if not, it means leakage. + private final Queue> queue = new ConcurrentLinkedQueue>(); + + private final ObjectPoolFactory factory; + + public DefaultObjectPool(ObjectPoolFactory factory, int size) { + if (factory == null) { + throw new NullPointerException("factory"); + } + this.factory = factory; + fill(size); + } + + private void fill(int size) { + for (int i = 0; i < size; i++) { + PooledObjectWrapper wrapper = createObject(); + queue.offer(wrapper); + } + } + + private PooledObjectWrapper createObject() { + T t = this.factory.create(); + return new PooledObjectWrapper(t); + } + + @Override + public PooledObject getObject() { + PooledObject object = queue.poll(); + if (object == null) { + // create dynamically ??? + return createObject(); + } + return object; + } + + + public void returnObject(PooledObject t) { + if (t == null) { + return; + } + factory.beforeReturn(t.getObject()); + queue.offer(t); + } + + private class PooledObjectWrapper implements PooledObject { + private final V value; + + public PooledObjectWrapper(V value) { + if (value == null) { + throw new NullPointerException("value must not be null"); + } + this.value = value; + } + + @Override + public V getObject() { + return value; + } + + @Override + public void returnObject() { + DefaultObjectPool.this.returnObject(this); + } + } + + +} diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/util/ObjectPool.java b/collector/src/main/java/com/navercorp/pinpoint/collector/util/ObjectPool.java index 349743404..1caeca74e 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/util/ObjectPool.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/util/ObjectPool.java @@ -16,50 +16,11 @@ package com.navercorp.pinpoint.collector.util; -import java.util.Queue; -import java.util.concurrent.ConcurrentLinkedQueue; - /** * @author emeroad */ -public class ObjectPool { - - // you don't need a blocking queue. There must be enough objects in a queue. if not, it means leakage. - private final Queue queue = new ConcurrentLinkedQueue(); - - private final ObjectPoolFactory factory; - - public ObjectPool(ObjectPoolFactory factory, int size) { - if (factory == null) { - throw new NullPointerException("factory"); - } - this.factory = factory; - fill(size); - } - - private void fill(int size) { - for (int i = 0; i < size; i++) { - T t = this.factory.create(); - queue.offer(t); - } - } - - public T getObject() { - T object = queue.poll(); - if (object == null) { - // create dynamically - return factory.create(); - } - return object; - } - - public void returnObject(T t) { - if (t == null) { - return; - } - factory.beforeReturn(t); - queue.offer(t); - } - +public interface ObjectPool { + PooledObject getObject(); +// void returnObject(PooledObject t); } diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/util/PooledObject.java b/collector/src/main/java/com/navercorp/pinpoint/collector/util/PooledObject.java new file mode 100644 index 000000000..cb8e78ec6 --- /dev/null +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/util/PooledObject.java @@ -0,0 +1,26 @@ +/* + * Copyright 2014 NAVER Corp. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.navercorp.pinpoint.collector.util; + +/** + * @author emeroad + */ +public interface PooledObject { + T getObject(); + + void returnObject(); +} diff --git a/collector/src/main/resources/applicationContext-collector.xml b/collector/src/main/resources/applicationContext-collector.xml index 6b6e1bc7b..0f6b4c328 100644 --- a/collector/src/main/resources/applicationContext-collector.xml +++ b/collector/src/main/resources/applicationContext-collector.xml @@ -87,24 +87,37 @@ - - + - - - - - + - - + + + + + + + + + + + + + + - - - - - + + + + + + + + + + + diff --git a/collector/src/test/java/com/navercorp/pinpoint/collector/receiver/udp/UDPReceiverTest.java b/collector/src/test/java/com/navercorp/pinpoint/collector/receiver/udp/UDPReceiverTest.java index d9f2dd4a0..583157ba3 100644 --- a/collector/src/test/java/com/navercorp/pinpoint/collector/receiver/udp/UDPReceiverTest.java +++ b/collector/src/test/java/com/navercorp/pinpoint/collector/receiver/udp/UDPReceiverTest.java @@ -23,7 +23,6 @@ import java.net.InetAddress; import java.net.InetSocketAddress; import java.net.SocketException; -import org.apache.thrift.TBase; import org.junit.Assert; import org.junit.Ignore; import org.junit.Test; @@ -31,7 +30,6 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import com.navercorp.pinpoint.collector.receiver.DataReceiver; -import com.navercorp.pinpoint.collector.receiver.DispatchHandler; /** * @author emeroad @@ -43,19 +41,12 @@ public class UDPReceiverTest { @Ignore public void startStop() { try { - DataReceiver receiver = new BaseUDPReceiver("test", new DispatchHandler() { + DataReceiver receiver = new UDPReceiver("test", new PacketHandlerFactory() { @Override - public void dispatchSendMessage(TBase tBase) { - } - - @Override - public TBase dispatchRequestMessage(TBase tBase) { - // TODO Auto-generated method stub + public PacketHandler createPacketHandler() { return null; } - }, "127.0.0.1", 10999, 1024, 1, 10); - } catch (Exception e) { e.printStackTrace(); Assert.fail(e.getMessage());