refactoring collector

- UDPReceiver
 - extract L4CheckFilter, NetworkAvailabilityCheckPacketFilter
This commit is contained in:
Woonduk Kang
2015-06-11 15:53:33 +09:00
parent 43492bda1e
commit ad552d4e2f
7 changed files with 234 additions and 50 deletions
@@ -27,7 +27,6 @@ import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.util.Assert;
import java.io.IOException;
import java.net.*;
/**
@@ -43,6 +42,8 @@ public class BaseUDPHandlerFactory<T extends DatagramPacket> implements PacketHa
private UDPReceiver receiver;
private final DispatchHandler dispatchHandler;
private TBaseFilter filter = TBaseFilter.CONTINUE_FILTER;
public BaseUDPHandlerFactory(DispatchHandler dispatchHandler) {
if (dispatchHandler == null) {
throw new NullPointerException("dispatchHandler must not be null");
@@ -54,6 +55,13 @@ public class BaseUDPHandlerFactory<T extends DatagramPacket> implements PacketHa
this.receiver = receiver;
}
public void setFilter(TBaseFilter filter) {
if (filter == null) {
throw new NullPointerException("filter must not be null");
}
this.filter = filter;
}
@Override
public PacketHandler<T> createPacketHandler() {
return new DispatchPacket();
@@ -75,19 +83,7 @@ public class BaseUDPHandlerFactory<T extends DatagramPacket> implements PacketHa
TBase<?, ?> tBase = null;
try {
tBase = deserializer.deserialize(packet.getData());
if (tBase instanceof L4Packet) {
if (logger.isDebugEnabled()) {
L4Packet l4Packet = (L4Packet) tBase;
logger.debug("udp l4 packet {}", l4Packet.getHeader());
}
return;
}
// Network port availability check packet
if (tBase instanceof NetworkAvailabilityCheckPacket) {
if (logger.isDebugEnabled()) {
logger.debug("received udp network availability check packet.");
}
responseOK(packet);
if (filter.filter(tBase, packet) == TBaseFilter.BREAK) {
return;
}
// dispatch signifies business logic execution
@@ -109,21 +105,6 @@ public class BaseUDPHandlerFactory<T extends DatagramPacket> implements PacketHa
}
}
}
private void responseOK(DatagramPacket packet) {
try {
byte[] okBytes = NetworkAvailabilityCheckPacket.DATA_OK;
DatagramPacket pongPacket = new DatagramPacket(okBytes, okBytes.length, packet.getSocketAddress());
receiver.getSocket().send(pongPacket);
} catch (IOException e) {
if (logger.isWarnEnabled()) {
logger.warn("pong error. SendSocketAddress:{} Cause:{}", packet.getSocketAddress(), e.getMessage(), e);
}
if (logger.isDebugEnabled()) {
logger.debug("packet dump hex:{}", PacketUtils.dumpDatagramPacket(packet));
}
}
}
}
}
@@ -44,6 +44,7 @@ public class ChunkedUDPPacketHandlerFactory<T extends DatagramPacket> implements
private UDPReceiver receiver;
private final DispatchHandler dispatchHandler;
private TBaseFilter filter = TBaseFilter.CONTINUE_FILTER;
public ChunkedUDPPacketHandlerFactory(DispatchHandler dispatchHandler) {
if (dispatchHandler == null) {
@@ -56,6 +57,12 @@ public class ChunkedUDPPacketHandlerFactory<T extends DatagramPacket> implements
this.receiver = receiver;
}
public void setFilter(TBaseFilter filter) {
if (filter == null) {
throw new NullPointerException("filter must not be null");
}
this.filter = filter;
}
@Override
public void afterPropertiesSet() throws Exception {
@@ -82,20 +89,8 @@ public class ChunkedUDPPacketHandlerFactory<T extends DatagramPacket> implements
}
for (TBase<?, ?> tBase : list) {
if (tBase instanceof L4Packet) {
if (logger.isDebugEnabled()) {
L4Packet l4Packet = (L4Packet) tBase;
logger.debug("udp l4 packet {}", l4Packet.getHeader());
}
continue;
}
// Network port availability check packet
if (tBase instanceof NetworkAvailabilityCheckPacket) {
if (logger.isDebugEnabled()) {
logger.debug("received udp network availability check packet.");
}
responseOK(packet);
continue;
if (filter.filter(tBase, packet) == TBaseFilter.BREAK) {
return;
}
// dispatch signifies business logic execution
dispatchHandler.dispatchSendMessage(tBase);
@@ -0,0 +1,44 @@
/*
* 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.thrift.io.L4Packet;
import org.apache.thrift.TBase;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.net.DatagramPacket;
/**
* @author emeroad
*/
public class L4PacketFilter implements TBaseFilter {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
@Override
public boolean filter(TBase<?, ?> tBase, DatagramPacket packet) {
if (tBase instanceof L4Packet) {
if (logger.isDebugEnabled()) {
L4Packet l4Packet = (L4Packet) tBase;
logger.debug("udp l4 packet {}", l4Packet.getHeader());
}
return BREAK;
}
return CONTINUE;
}
}
@@ -0,0 +1,73 @@
/*
* 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.PacketUtils;
import com.navercorp.pinpoint.thrift.io.NetworkAvailabilityCheckPacket;
import org.apache.thrift.TBase;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.net.DatagramPacket;
import java.net.DatagramSocket;
import java.net.SocketException;
/**
* @author emeroad
*/
public class NetworkAvailabilityCheckPacketFilter implements TBaseFilter {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
private final DatagramSocket socket;
public NetworkAvailabilityCheckPacketFilter() {
try {
this.socket = new DatagramSocket();
} catch (SocketException ex) {
throw new RuntimeException("socket create fail. error:" + ex.getMessage(), ex);
}
}
@Override
public boolean filter(TBase<?, ?> tBase, DatagramPacket packet) {
// Network port availability check packet
if (tBase instanceof NetworkAvailabilityCheckPacket) {
if (logger.isInfoEnabled()) {
logger.info("received udp network availability check packet.");
}
responseOK(packet);
return BREAK;
}
return CONTINUE;
}
private void responseOK(DatagramPacket packet) {
try {
byte[] okBytes = NetworkAvailabilityCheckPacket.DATA_OK;
DatagramPacket pongPacket = new DatagramPacket(okBytes, okBytes.length, packet.getSocketAddress());
socket.send(pongPacket);
} catch (IOException e) {
if (logger.isWarnEnabled()) {
logger.warn("pong error. SendSocketAddress:{} Cause:{}", packet.getSocketAddress(), e.getMessage(), e);
}
if (logger.isDebugEnabled()) {
logger.debug("packet dump hex:{}", PacketUtils.dumpDatagramPacket(packet));
}
}
}
}
@@ -0,0 +1,39 @@
/*
* 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 org.apache.thrift.TBase;
import java.net.DatagramPacket;
/**
* @author emeroad
*/
public interface TBaseFilter {
boolean CONTINUE = true;
boolean BREAK = false;
boolean filter(TBase<?, ?> tBase, DatagramPacket packet);
public static final TBaseFilter CONTINUE_FILTER = new TBaseFilter() {
@Override
public boolean filter(TBase<?, ?> tBase, DatagramPacket packet) {
return CONTINUE;
}
};
}
@@ -0,0 +1,48 @@
/*
* 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 org.apache.thrift.TBase;
import java.net.DatagramPacket;
import java.util.ArrayList;
import java.util.List;
/**
* @author emeroad
*/
public class TBaseFilterChain implements TBaseFilter {
private final List<TBaseFilter> filterChain = new ArrayList<TBaseFilter>();
public void addTBaseFilter(TBaseFilter tBaseFilter) {
if (tBaseFilter == null) {
throw new NullPointerException("tBaseFilter must not be null");
}
this.filterChain.add(tBaseFilter);
}
@Override
public boolean filter(TBase<?, ?> tBase, DatagramPacket packet) {
for (TBaseFilter tBaseFilter : filterChain) {
if (tBaseFilter.filter(tBase, packet) == TBaseFilter.BREAK) {
return BREAK;
}
}
return TBaseFilter.CONTINUE;
}
}
@@ -142,7 +142,7 @@ public class UDPReceiver implements DataReceiver {
if (pooledPacket == null) {
continue;
}
DatagramPacket packet = pooledPacket.getObject();
final DatagramPacket packet = pooledPacket.getObject();
if (packet.getLength() == 0) {
if (debugEnabled) {
logger.debug("length is 0 ip:{}, port:{}", packet.getAddress(), packet.getPort());
@@ -156,12 +156,7 @@ public class UDPReceiver implements DataReceiver {
Runnable dispatchTask = wrapDispatchTask(pooledPacket);
worker.execute(dispatchTask);
} catch (RejectedExecutionException ree) {
rejectedCounter.inc();
final int error = rejectedExecutionCount.incrementAndGet();
final int mod = 100;
if ((error % mod) == 0) {
logger.warn("RejectedExecutionCount={}", error);
}
handleRejectedExecutionException(ree);
}
}
if (logger.isInfoEnabled()) {
@@ -169,6 +164,15 @@ public class UDPReceiver implements DataReceiver {
}
}
private void handleRejectedExecutionException(RejectedExecutionException ree) {
rejectedCounter.inc();
final int error = rejectedExecutionCount.incrementAndGet();
final int mod = 100;
if ((error % mod) == 0) {
logger.warn("RejectedExecutionCount={}", error);
}
}
private Runnable wrapDispatchTask(PooledObject<DatagramPacket> pooledPacket) {
PacketHandler<DatagramPacket> dispatchPacket = packetHandlerFactory.createPacketHandler();
PooledPacketWrap pooledPacketWrap = new PooledPacketWrap(dispatchPacket, pooledPacket);