mirror of
https://github.com/wahyd4/pinpoint.git
synced 2026-08-16 08:16:15 +10:00
refactoring collector
- UDPReceiver - ObjectPool
This commit is contained in:
+37
-29
@@ -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<T extends DatagramPacket> implements PacketHandlerFactory<T>, InitializingBean {
|
||||
|
||||
private final Logger logger = LoggerFactory.getLogger(this.getClass());
|
||||
|
||||
private DeserializerFactory<HeaderTBaseDeserializer> deserializerFactory = new ThreadLocalHeaderTBaseDeserializerFactory<HeaderTBaseDeserializer>(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<T> 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<T> {
|
||||
|
||||
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());
|
||||
+36
-28
@@ -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<T extends DatagramPacket> implements PacketHandlerFactory<T>, InitializingBean {
|
||||
|
||||
private final Logger logger = LoggerFactory.getLogger(this.getClass());
|
||||
|
||||
private final DeserializerFactory<ChunkHeaderTBaseDeserializer> deserializerFactory = new ThreadLocalHeaderTBaseDeserializerFactory<ChunkHeaderTBaseDeserializer>(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<T> 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<T> {
|
||||
|
||||
private DispatchPacket() {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
Timer.Context time = receiver.getTimer().time();
|
||||
|
||||
public void receive(T packet) {
|
||||
final ChunkHeaderTBaseDeserializer deserializer = deserializerFactory.createDeserializer();
|
||||
try {
|
||||
List<TBase<?, ?>> 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());
|
||||
+24
@@ -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<T> {
|
||||
void receive(T packet);
|
||||
}
|
||||
+24
@@ -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<T> {
|
||||
PacketHandler<T> createPacketHandler();
|
||||
}
|
||||
+50
@@ -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<DatagramPacket> packetHandler;
|
||||
private final PooledObject<DatagramPacket> pooledObject;
|
||||
|
||||
public PooledPacketWrap(PacketHandler<DatagramPacket> packetHandler, PooledObject<DatagramPacket> 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();
|
||||
}
|
||||
}
|
||||
}
|
||||
+26
-26
@@ -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<T extends DatagramPacket> implements PacketHandlerFactory<T>, InitializingBean {
|
||||
|
||||
private final Logger logger = LoggerFactory.getLogger(this.getClass());
|
||||
|
||||
private DeserializerFactory<HeaderTBaseDeserializer> deserializerFactory = new ThreadLocalHeaderTBaseDeserializerFactory<HeaderTBaseDeserializer>(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<T> 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<T> {
|
||||
|
||||
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();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
+37
-38
@@ -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<DatagramPacket> 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<DatagramPacket> 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<DatagramPacket>(new DatagramPacketFactory(), packetPoolSize);
|
||||
this.datagramPacketPool = new DefaultObjectPool<DatagramPacket>(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<DatagramPacket> 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<DatagramPacket> pooledPacket) {
|
||||
PacketHandler<DatagramPacket> dispatchPacket = packetHandlerFactory.createPacketHandler();
|
||||
PooledPacketWrap pooledPacketWrap = new PooledPacketWrap(dispatchPacket, pooledPacket);
|
||||
return new TimingWrap(this.timer, pooledPacketWrap);
|
||||
}
|
||||
|
||||
private PooledObject<DatagramPacket> read0() {
|
||||
boolean success = false;
|
||||
DatagramPacket packet = datagramPacketPool.getObject();
|
||||
if (packet == null) {
|
||||
PooledObject<DatagramPacket> 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<DatagramPacket> getDatagramPacketPool() {
|
||||
return datagramPacketPool;
|
||||
}
|
||||
|
||||
public DatagramSocket getSocket() {
|
||||
return socket;
|
||||
}
|
||||
@@ -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<ByteBuffer> {
|
||||
|
||||
private static final int AcceptedSize = 65507;
|
||||
|
||||
@Override
|
||||
public ByteBuffer create() {
|
||||
return ByteBuffer.allocate(AcceptedSize);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void beforeReturn(ByteBuffer packet) {
|
||||
packet.clear();
|
||||
}
|
||||
}
|
||||
@@ -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<T> implements ObjectPool<T> {
|
||||
|
||||
// you don't need a blocking queue. There must be enough objects in a queue. if not, it means leakage.
|
||||
private final Queue<PooledObject<T>> queue = new ConcurrentLinkedQueue<PooledObject<T>>();
|
||||
|
||||
private final ObjectPoolFactory<T> factory;
|
||||
|
||||
public DefaultObjectPool(ObjectPoolFactory<T> 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<T> wrapper = createObject();
|
||||
queue.offer(wrapper);
|
||||
}
|
||||
}
|
||||
|
||||
private PooledObjectWrapper<T> createObject() {
|
||||
T t = this.factory.create();
|
||||
return new PooledObjectWrapper<T>(t);
|
||||
}
|
||||
|
||||
@Override
|
||||
public PooledObject<T> getObject() {
|
||||
PooledObject<T> object = queue.poll();
|
||||
if (object == null) {
|
||||
// create dynamically ???
|
||||
return createObject();
|
||||
}
|
||||
return object;
|
||||
}
|
||||
|
||||
|
||||
public void returnObject(PooledObject<T> t) {
|
||||
if (t == null) {
|
||||
return;
|
||||
}
|
||||
factory.beforeReturn(t.getObject());
|
||||
queue.offer(t);
|
||||
}
|
||||
|
||||
private class PooledObjectWrapper<V extends T > implements PooledObject<T> {
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
@@ -16,50 +16,11 @@
|
||||
|
||||
package com.navercorp.pinpoint.collector.util;
|
||||
|
||||
import java.util.Queue;
|
||||
import java.util.concurrent.ConcurrentLinkedQueue;
|
||||
|
||||
/**
|
||||
* @author emeroad
|
||||
*/
|
||||
public class ObjectPool<T> {
|
||||
|
||||
// you don't need a blocking queue. There must be enough objects in a queue. if not, it means leakage.
|
||||
private final Queue<T> queue = new ConcurrentLinkedQueue<T>();
|
||||
|
||||
private final ObjectPoolFactory<T> factory;
|
||||
|
||||
public ObjectPool(ObjectPoolFactory<T> 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<T> {
|
||||
PooledObject<T> getObject();
|
||||
|
||||
// void returnObject(PooledObject<T> t);
|
||||
}
|
||||
|
||||
@@ -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> {
|
||||
T getObject();
|
||||
|
||||
void returnObject();
|
||||
}
|
||||
@@ -87,24 +87,37 @@
|
||||
<constructor-arg type="com.navercorp.pinpoint.collector.cluster.zookeeper.ZookeeperClusterService" ref="clusterService"/>
|
||||
</bean>
|
||||
|
||||
<bean id="udpSpanReceiver" class="com.navercorp.pinpoint.collector.receiver.udp.BaseUDPReceiver">
|
||||
<constructor-arg value="Pinpoint-UDP-Span"/>
|
||||
<bean id="udpSpanBasePacketHandler" class="com.navercorp.pinpoint.collector.receiver.udp.BaseUDPHandlerFactory">
|
||||
<constructor-arg type="com.navercorp.pinpoint.collector.receiver.DispatchHandler" ref="udpSpanDispatchHandler"/>
|
||||
<constructor-arg value="#{collectorConfiguration.udpSpanListenIp}"/>
|
||||
<constructor-arg value="#{collectorConfiguration.udpSpanListenPort}"/>
|
||||
<constructor-arg value="#{collectorConfiguration.udpSpanSocketReceiveBufferSize}"/>
|
||||
<constructor-arg value="#{collectorConfiguration.udpSpanWorkerThread}"/>
|
||||
<constructor-arg value="#{collectorConfiguration.udpSpanWorkerQueueSize}"/>
|
||||
<property name="receiver" ref="udpSpanReceiver"/>
|
||||
</bean>
|
||||
|
||||
<bean id="udpStatReceiver" class="com.navercorp.pinpoint.collector.receiver.udp.BaseUDPReceiver">
|
||||
<constructor-arg value="Pinpoint-UDP-Stat"/>
|
||||
|
||||
<bean id="udpSpanReceiver" class="com.navercorp.pinpoint.collector.receiver.udp.UDPReceiver">
|
||||
<constructor-arg index="0" value="Pinpoint-UDP-Span"/>
|
||||
<constructor-arg index="1" ref="udpSpanBasePacketHandler"/>
|
||||
<constructor-arg index="2" value="#{collectorConfiguration.udpSpanListenIp}"/>
|
||||
<constructor-arg index="3" value="#{collectorConfiguration.udpSpanListenPort}"/>
|
||||
<constructor-arg index="4" value="#{collectorConfiguration.udpSpanSocketReceiveBufferSize}"/>
|
||||
<constructor-arg index="5" value="#{collectorConfiguration.udpSpanWorkerThread}"/>
|
||||
<constructor-arg index="6" value="#{collectorConfiguration.udpSpanWorkerQueueSize}"/>
|
||||
|
||||
</bean>
|
||||
|
||||
|
||||
<bean id="udpStatBasePacketHandler" class="com.navercorp.pinpoint.collector.receiver.udp.BaseUDPHandlerFactory">
|
||||
<constructor-arg type="com.navercorp.pinpoint.collector.receiver.DispatchHandler" ref="udpDispatchHandler"/>
|
||||
<constructor-arg value="#{collectorConfiguration.udpStatListenIp}"/>
|
||||
<constructor-arg value="#{collectorConfiguration.udpStatListenPort}"/>
|
||||
<constructor-arg value="#{collectorConfiguration.udpStatSocketReceiveBufferSize}"/>
|
||||
<constructor-arg value="#{collectorConfiguration.udpStatWorkerThread}"/>
|
||||
<constructor-arg value="#{collectorConfiguration.udpStatWorkerQueueSize}"/>
|
||||
<property name="receiver" ref="udpStatReceiver"/>
|
||||
</bean>
|
||||
|
||||
<bean id="udpStatReceiver" class="com.navercorp.pinpoint.collector.receiver.udp.UDPReceiver">
|
||||
<constructor-arg index="0" value="Pinpoint-UDP-Stat"/>
|
||||
<constructor-arg index="1" ref="udpStatBasePacketHandler"/>
|
||||
<constructor-arg index="2" value="#{collectorConfiguration.udpStatListenIp}"/>
|
||||
<constructor-arg index="3" value="#{collectorConfiguration.udpStatListenPort}"/>
|
||||
<constructor-arg index="4" value="#{collectorConfiguration.udpStatSocketReceiveBufferSize}"/>
|
||||
<constructor-arg index="5" value="#{collectorConfiguration.udpStatWorkerThread}"/>
|
||||
<constructor-arg index="6" value="#{collectorConfiguration.udpStatWorkerQueueSize}"/>
|
||||
</bean>
|
||||
|
||||
<bean id="jsonObjectMapper" class="org.codehaus.jackson.map.ObjectMapper">
|
||||
|
||||
+2
-11
@@ -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());
|
||||
|
||||
Reference in New Issue
Block a user