From bd03dfadae2d7be302ef26efba92f63cb4691f71 Mon Sep 17 00:00:00 2001 From: Jaehong Kim Date: Wed, 5 Nov 2014 18:59:28 +0900 Subject: [PATCH] #28 change class --- .../collector/cluster/ClusterPointRouter.java | 2 +- .../collector/receiver/tcp/TCPReceiver.java | 2 +- .../profiler/receiver/CommandDispatcher.java | 8 +-- .../io/ByteArrayOutputStreamTransport.java | 2 +- .../ChunkHeaderBufferedTBaseSerializer.java | 49 ++++++++----------- ...erBufferedTBaseSerializerFlushHandler.java | 6 ++- .../io/ChunkHeaderTBaseDeserializer.java | 15 +----- .../thrift/io/DefaultTBaseLocator.java | 2 + .../io/HeaderTBaseSerializerFactory.java | 2 +- .../pinpoint/thrift/io/SerializerFactory.java | 4 +- ...readLocalHeaderTBaseSerializerFactory.java | 15 +++--- ...hunkHeaderBufferedTBaseSerializerTest.java | 6 ++- 12 files changed, 50 insertions(+), 63 deletions(-) diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/ClusterPointRouter.java b/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/ClusterPointRouter.java index e8996d21d..771ff1472 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/ClusterPointRouter.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/ClusterPointRouter.java @@ -42,7 +42,7 @@ public class ClusterPointRouter { private final WebClusterPoint webClusterPoint; @Autowired - private SerializerFactory commandSerializerFactory; + private SerializerFactory commandSerializerFactory; @Autowired private DeserializerFactory commandDeserializerFactory; diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/TCPReceiver.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/TCPReceiver.java index 6b94954c1..99c8e6c64 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/TCPReceiver.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/TCPReceiver.java @@ -65,7 +65,7 @@ public class TCPReceiver { private final ThreadPoolExecutor worker = ExecutorFactory.newFixedThreadPool(threadSize, workerQueueSize, THREAD_FACTORY); - private final SerializerFactory serializerFactory = new ThreadLocalHeaderTBaseSerializerFactory(new HeaderTBaseSerializerFactory(true, HeaderTBaseSerializerFactory.DEFAULT_UDP_STREAM_MAX_SIZE)); + private final SerializerFactory serializerFactory = new ThreadLocalHeaderTBaseSerializerFactory(new HeaderTBaseSerializerFactory(true, HeaderTBaseSerializerFactory.DEFAULT_UDP_STREAM_MAX_SIZE)); private final DeserializerFactory deserializerFactory = new ThreadLocalHeaderTBaseDeserializerFactory(new HeaderTBaseDeserializerFactory()); diff --git a/profiler/src/main/java/com/navercorp/pinpoint/profiler/receiver/CommandDispatcher.java b/profiler/src/main/java/com/navercorp/pinpoint/profiler/receiver/CommandDispatcher.java index 9bcf37445..ba8ce6fc9 100644 --- a/profiler/src/main/java/com/navercorp/pinpoint/profiler/receiver/CommandDispatcher.java +++ b/profiler/src/main/java/com/navercorp/pinpoint/profiler/receiver/CommandDispatcher.java @@ -38,7 +38,7 @@ public class CommandDispatcher implements MessageListener { private final ProfilerCommandServiceLocator locator; - private final SerializerFactory serializerFactory; + private final SerializerFactory serializerFactory; private final DeserializerFactory deserializerFactory; public CommandDispatcher(Builder builder) { @@ -48,7 +48,7 @@ public class CommandDispatcher implements MessageListener { } this.locator = registry; - SerializerFactory serializerFactory = new HeaderTBaseSerializerFactory(true, builder.serializationMaxSize, builder.protocolFactory, builder.commandTbaseLocator); + SerializerFactory serializerFactory = new HeaderTBaseSerializerFactory(true, builder.serializationMaxSize, builder.protocolFactory, builder.commandTbaseLocator); this.serializerFactory = wrappedThreadLocalSerializerFactory(serializerFactory); AssertUtils.assertNotNull(this.serializerFactory); @@ -57,8 +57,8 @@ public class CommandDispatcher implements MessageListener { AssertUtils.assertNotNull(this.deserializerFactory); } - private SerializerFactory wrappedThreadLocalSerializerFactory(SerializerFactory serializerFactory) { - return new ThreadLocalHeaderTBaseSerializerFactory(serializerFactory); + private SerializerFactory wrappedThreadLocalSerializerFactory(SerializerFactory serializerFactory) { + return new ThreadLocalHeaderTBaseSerializerFactory(serializerFactory); } private DeserializerFactory wrappedThreadLocalDeserializerFactory(DeserializerFactory deserializerFactory) { diff --git a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ByteArrayOutputStreamTransport.java b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ByteArrayOutputStreamTransport.java index e62b40f21..59c4d02b2 100644 --- a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ByteArrayOutputStreamTransport.java +++ b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ByteArrayOutputStreamTransport.java @@ -8,7 +8,7 @@ import org.apache.thrift.transport.TTransportException; /** * ByteArrayOutputStreamTransport - * - unsupported read operation + * - write only * * @author jaehong.kim */ diff --git a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderBufferedTBaseSerializer.java b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderBufferedTBaseSerializer.java index f57aab1c7..f5e24ba69 100644 --- a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderBufferedTBaseSerializer.java +++ b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderBufferedTBaseSerializer.java @@ -12,6 +12,12 @@ import com.nhn.pinpoint.thrift.dto.TSpan; import com.nhn.pinpoint.thrift.dto.TSpanChunk; import com.nhn.pinpoint.thrift.dto.TSpanEvent; +/** + * ChunkHeaderBufferedTBaseSerializer + * - need flush handler + * + * @author jaehong.kim + */ public class ChunkHeaderBufferedTBaseSerializer { private static final String FIELD_NAME_SPAN_EVENT_LIST = "spanEventList"; @@ -47,6 +53,7 @@ public class ChunkHeaderBufferedTBaseSerializer { } } + // TSpanChunk = TSpanChunk + TSpanChunk private void addTSpanChunk(TBase base) throws TException { final TSpanChunk chunk = (TSpanChunk) base; if (chunk.getSpanEventList() == null) { @@ -67,6 +74,7 @@ public class ChunkHeaderBufferedTBaseSerializer { } } + // TSpan = TSpan + TSpanChunk private void addTSpan(TBase base) throws TException { final TSpan span = (TSpan) base; if (span.getSpanEventList() == null) { @@ -88,12 +96,13 @@ public class ChunkHeaderBufferedTBaseSerializer { } } + // write chunk header + header + body private void write(final TBase base, final String fieldName, final List list) throws TException { final ReplaceListCompactProtocol protocol = new ReplaceListCompactProtocol(new ByteArrayOutputStreamTransport(out)); // write chunk header writeChunkHeader(protocol); - + // write header writeHeader(protocol, locator.headerLookup(base)); if (list != null && list.size() > 0) { @@ -101,12 +110,13 @@ public class ChunkHeaderBufferedTBaseSerializer { } base.write(protocol); - + if (isOverflow()) { flush(); } } + // write chunk header + header + body private void write(final TBase base) throws TException { final TCompactProtocol protocol = new TCompactProtocol(new ByteArrayOutputStreamTransport(out)); @@ -117,7 +127,7 @@ public class ChunkHeaderBufferedTBaseSerializer { writeHeader(protocol, locator.headerLookup(base)); base.write(protocol); - + if (isOverflow()) { flush(); } @@ -140,7 +150,6 @@ public class ChunkHeaderBufferedTBaseSerializer { private void writeHeader(final TProtocol protocol, final Header header) throws TException { protocol.writeByte(header.getSignature()); protocol.writeByte(header.getVersion()); - // 프로토콜 변경에 관계 없이 고정 사이즈의 데이터로 인코딩 하도록 변경. short type = header.getType(); protocol.writeByte(BytesUtils.writeShort1(type)); protocol.writeByte(BytesUtils.writeShort2(type)); @@ -166,15 +175,13 @@ public class ChunkHeaderBufferedTBaseSerializer { } public String toString() { - - return toStringBinary(out.toByteArray(), 0, out.size()); -// StringBuilder sb = new StringBuilder(); -// sb.append("{"); -// sb.append("bufferSize=").append(out.size()).append(", "); -// sb.append("flushSize=").append(flushSize); -// sb.append("}"); -// -// return sb.toString(); + StringBuilder sb = new StringBuilder(); + sb.append("{"); + sb.append("bufferSize=").append(out.size()).append(", "); + sb.append("flushSize=").append(flushSize); + sb.append("}"); + + return sb.toString(); } TSpanChunk toSpanChunk(TSpan span) { @@ -199,20 +206,4 @@ public class ChunkHeaderBufferedTBaseSerializer { return spanChunk; } - - - - static String toStringBinary(final byte[] b, int off, int len) { - StringBuilder result = new StringBuilder(); - for (int i = off; i < off + len; ++i) { - int ch = b[i] & 0xFF; - if ((ch >= '0' && ch <= '9') || (ch >= 'A' && ch <= 'Z') || (ch >= 'a' && ch <= 'z') || " `~!@#$%^&*()-_=+[]{}|;:'\",.<>/?".indexOf(ch) >= 0) { - result.append((char) ch); - } else { - result.append(String.format("\\x%02X", ch)); - } - } - return result.toString(); - } - } \ No newline at end of file diff --git a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderBufferedTBaseSerializerFlushHandler.java b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderBufferedTBaseSerializerFlushHandler.java index 4442bca63..d613ef318 100644 --- a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderBufferedTBaseSerializerFlushHandler.java +++ b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderBufferedTBaseSerializerFlushHandler.java @@ -1,6 +1,10 @@ package com.nhn.pinpoint.thrift.io; +/** + * + * @author jaehong.kim + */ public interface ChunkHeaderBufferedTBaseSerializerFlushHandler { void handle(byte[] buffer, int offset, int length); -} +} \ No newline at end of file diff --git a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderTBaseDeserializer.java b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderTBaseDeserializer.java index 693ff0bda..c0f5eb8ea 100644 --- a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderTBaseDeserializer.java +++ b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderTBaseDeserializer.java @@ -14,30 +14,18 @@ public class ChunkHeaderTBaseDeserializer { private final TMemoryInputTransport trans; private final TBaseLocator locator; - /** - * Create a new TDeserializer. It will use the TProtocol specified by the factory that is passed in. - * - * @param protocolFactory - * Factory to create a protocol - */ ChunkHeaderTBaseDeserializer(TProtocolFactory protocolFactory, TBaseLocator locator) { this.trans = new TMemoryInputTransport(); this.protocol = protocolFactory.getProtocol(trans); this.locator = locator; } - /** - * Deserialize the Thrift object from a byte array. - * - * @param bytes - * The array to read from - */ public List> deserialize(byte[] bytes, int offset, int length) throws TException { List> list = new ArrayList>(); try { trans.reset(bytes, offset, length); - Header header = readHeader(); + final Header header = readHeader(); if (locator.isChunkHeader(header.getType())) { TBase base = null; @@ -92,7 +80,6 @@ public class ChunkHeaderTBaseDeserializer { final byte signature = protocol.readByte(); final byte version = protocol.readByte(); - // 프로토콜 변경에 관계 없이 고정 사이즈의 데이터로 인코딩 하도록 변경. final byte type1 = protocol.readByte(); final byte type2 = protocol.readByte(); final short type = bytesToShort(type1, type2); diff --git a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/DefaultTBaseLocator.java b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/DefaultTBaseLocator.java index 1aaa305e3..538acd351 100644 --- a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/DefaultTBaseLocator.java +++ b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/DefaultTBaseLocator.java @@ -18,6 +18,8 @@ import com.nhn.pinpoint.thrift.dto.TStringMetaData; * @author koo.taejin * @author netspider * @author hyungil.jeong + * @author jaehong.kim + * - add CHUNK_HEADER */ class DefaultTBaseLocator implements TBaseLocator { diff --git a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/HeaderTBaseSerializerFactory.java b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/HeaderTBaseSerializerFactory.java index bf9b35946..8c61b394b 100644 --- a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/HeaderTBaseSerializerFactory.java +++ b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/HeaderTBaseSerializerFactory.java @@ -8,7 +8,7 @@ import java.io.ByteArrayOutputStream; /** * @author koo.taejin */ -public final class HeaderTBaseSerializerFactory implements SerializerFactory { +public final class HeaderTBaseSerializerFactory implements SerializerFactory { private static final TBaseLocator DEFAULT_TBASE_LOCATOR = new DefaultTBaseLocator(); diff --git a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/SerializerFactory.java b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/SerializerFactory.java index d9c53da13..14a077c97 100644 --- a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/SerializerFactory.java +++ b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/SerializerFactory.java @@ -3,6 +3,6 @@ package com.nhn.pinpoint.thrift.io; /** * @author emeroad */ -public interface SerializerFactory { - HeaderTBaseSerializer createSerializer(); +public interface SerializerFactory { + E createSerializer(); } diff --git a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ThreadLocalHeaderTBaseSerializerFactory.java b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ThreadLocalHeaderTBaseSerializerFactory.java index ce3ccc30d..77187faee 100644 --- a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ThreadLocalHeaderTBaseSerializerFactory.java +++ b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/ThreadLocalHeaderTBaseSerializerFactory.java @@ -2,28 +2,29 @@ package com.nhn.pinpoint.thrift.io; /** * @author emeroad + * @author jaehong.kim + * - change to generic type */ -public class ThreadLocalHeaderTBaseSerializerFactory implements SerializerFactory { +public class ThreadLocalHeaderTBaseSerializerFactory implements SerializerFactory { - private final ThreadLocal cache = new ThreadLocal() { + private final ThreadLocal cache = new ThreadLocal() { @Override - protected HeaderTBaseSerializer initialValue() { + protected E initialValue() { return factory.createSerializer(); } }; - private final SerializerFactory factory; + private final SerializerFactory factory; - public ThreadLocalHeaderTBaseSerializerFactory(SerializerFactory factory) { + public ThreadLocalHeaderTBaseSerializerFactory(SerializerFactory factory) { if (factory == null) { throw new NullPointerException("factory must not be null"); } this.factory = factory; } - @Override - public HeaderTBaseSerializer createSerializer() { + public E createSerializer() { return cache.get(); } } diff --git a/thrift/src/test/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderBufferedTBaseSerializerTest.java b/thrift/src/test/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderBufferedTBaseSerializerTest.java index 3c40793b6..7e70096b1 100644 --- a/thrift/src/test/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderBufferedTBaseSerializerTest.java +++ b/thrift/src/test/java/com/navercorp/pinpoint/thrift/io/ChunkHeaderBufferedTBaseSerializerTest.java @@ -3,7 +3,6 @@ package com.nhn.pinpoint.thrift.io; import static org.junit.Assert.*; import java.util.Arrays; -import java.util.concurrent.TimeUnit; import org.apache.thrift.TException; import org.junit.Test; @@ -16,7 +15,6 @@ public class ChunkHeaderBufferedTBaseSerializerTest { @Test public void add() throws TException { - System.out.println("add"); ChunkHeaderBufferedTBaseSerializer serializer = new ChunkHeaderBufferedTBaseSerializer(1024); serializer.setFlushHandler(new ChunkHeaderBufferedTBaseSerializerFlushHandler() { @@ -27,24 +25,28 @@ public class ChunkHeaderBufferedTBaseSerializerTest { } }); + // add and flush flush = false; TSpanChunk chunk = new TSpanMockBuilder().buildChunk(1, 1024); serializer.add(chunk); System.out.println(serializer); assertTrue(flush); + // add and flush * 3 flush = false; chunk = new TSpanMockBuilder().buildChunk(3, 1024); serializer.add(chunk); System.out.println(serializer); assertTrue(flush); + // add flush = false; chunk = new TSpanMockBuilder().buildChunk(3, 10); serializer.add(chunk); System.out.println(serializer); assertFalse(flush); + // flush serializer.flush(); assertTrue(flush); }