From 54c8dbf7671dd8e51871f68ad64b8f6ab1da6291 Mon Sep 17 00:00:00 2001 From: wusheng Date: Fri, 15 Jul 2016 10:39:48 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E6=94=B9=E5=A4=A7=E9=87=8F=E7=9A=84?= =?UTF-8?q?=E6=96=B9=E6=B3=95=E4=BF=AE=E6=94=B9=E6=8C=87=E5=AF=BC=EF=BC=8C?= =?UTF-8?q?=E4=BF=AE=E6=AD=A3=E9=94=99=E8=AF=AF=E7=9A=84=E6=96=B9=E6=B3=95?= =?UTF-8?q?=E9=80=BB=E8=BE=91=EF=BC=8C=E6=8F=90=E9=AB=98=E6=96=B9=E6=B3=95?= =?UTF-8?q?=E7=9A=84=E5=8F=AF=E8=AF=BB=E6=80=A7=E5=92=8C=E6=95=88=E7=8E=87?= =?UTF-8?q?=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../protocol/TransportPackager.java | 2 +- ...BalePlucker.java => BufferFileReader.java} | 30 +++++++++++++++---- .../peresistent/PersistenceThread.java | 30 ++++--------------- 3 files changed, 32 insertions(+), 30 deletions(-) rename skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/{BufferBalePlucker.java => BufferFileReader.java} (87%) diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/TransportPackager.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/TransportPackager.java index c3eac035f..54c60035d 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/TransportPackager.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/TransportPackager.java @@ -16,7 +16,7 @@ public class TransportPackager { return dataPackage; } - public static List unpackDataBody(byte[] dataPackage) { + public static List unpackDataBody(byte[] dataPackage) { List serializeData = null; try { serializeData = new ArrayList(); diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/BufferBalePlucker.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/BufferFileReader.java similarity index 87% rename from skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/BufferBalePlucker.java rename to skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/BufferFileReader.java index e5a378da6..17f5a5f01 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/BufferBalePlucker.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/BufferFileReader.java @@ -1,5 +1,6 @@ package com.ai.cloud.skywalking.reciever.peresistent; +import com.ai.cloud.skywalking.protocol.common.AbstractDataSerializable; import com.ai.cloud.skywalking.protocol.util.IntegerAssist; import org.apache.commons.io.FileUtils; import org.apache.logging.log4j.LogManager; @@ -9,8 +10,9 @@ import java.io.File; import java.io.FileInputStream; import java.io.IOException; import java.util.Arrays; +import java.util.List; -public class BufferBalePlucker { +public class BufferFileReader { private File bufferFile; private FileInputStream bufferInputStream; private int currentOffset; @@ -19,10 +21,10 @@ public class BufferBalePlucker { private static final byte[] EOF_BALE_ARRAY = "EOF".getBytes(); private int remainderLength = 0; private byte[] remainderByte = null; - private Logger logger = LogManager.getLogger(BufferBalePlucker.class); + private Logger logger = LogManager.getLogger(BufferFileReader.class); private PluckerStatus status = PluckerStatus.INITIAL; - public BufferBalePlucker(File bufferFile, int currentOffset) { + public BufferFileReader(File bufferFile, int currentOffset) { this.bufferFile = bufferFile; this.currentOffset = currentOffset; @@ -39,7 +41,25 @@ public class BufferBalePlucker { } - public boolean hasNextBufferBale() { + /** + * TODO: 确保读到下一个List,并进行缓存 + * + * 按照如下格式读取: + * 1.头4位长度 + * 2.根据长度读取正文 + * 2.1. 正文包含多个ISerializable。每个块包含每个ISerializable的长度和ISerializable的正文 + * 3.长度外,读取4位,为分隔符(不可见字符) + * + * 封装方法要求: + * 1.封装读取指定长度块的方法,读取完成则返回。读取长度不足,则缓存并等待。 + * + * 异常处理: + * 1.长度外,读取4位,不是分隔符: 则启动异常跳位处理,直到读取到分隔符位置 + * + * + * @return + */ + public boolean hasNext() { if (status == PluckerStatus.SUSPEND) { try { @@ -70,7 +90,7 @@ public class BufferBalePlucker { return hasNextBufferBale; } - public byte[] pluck() throws IOException { + public List next() throws IOException { int packageLength = unpackBaleLength(); byte[] dataPackage = unpackDataContext(packageLength); diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThread.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThread.java index 7efbe2e60..55159a26d 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThread.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThread.java @@ -45,14 +45,14 @@ public class PersistenceThread extends Thread { int offset = acquireOffset(); - BufferBalePlucker bufferBalePlucker = new BufferBalePlucker(bufferFile, offset); - while (bufferBalePlucker.hasNextBufferBale()) { + BufferFileReader bufferReader = new BufferFileReader(bufferFile, offset); + while (bufferReader.hasNext()) { try { - Map> spans = unSerializeSpans(bufferBalePlucker.pluck()); - System.out.println(spans.size()); + List serializableDataList = bufferReader.next(); + System.out.println(serializableDataList.size()); //handleSpans(spans); } catch (ConvertFailedException e) { - bufferBalePlucker.skipToNextBufferBale(); + bufferReader.skipToNextBufferBale(); } catch (IOException e) { logger.error("The data file I/O exception.", e); ServerHealthCollector.getCurrentHeathReading(null) @@ -61,7 +61,7 @@ public class PersistenceThread extends Thread { } try { - bufferBalePlucker.close(); + bufferReader.close(); } catch (IOException e) { logger.error("The data file I/O exception.", e); ServerHealthCollector.getCurrentHeathReading(null).updateData(ServerHeathReading.ERROR, e.getMessage()); @@ -93,24 +93,6 @@ public class PersistenceThread extends Thread { } } - private Map> unSerializeSpans(byte[] byteData) - throws ConvertFailedException { - Map> spans = new HashMap>(); - if (byteData == null || byteData.length == 0) { - return spans; - } - List serializeData = TransportPackager.unpackDataBody(byteData); - for (byte[] dataBytes : serializeData) { - AbstractDataSerializable abstractDataSerializable = SerializedFactory.unSerialize(dataBytes); - if (spans.get(abstractDataSerializable.getDataType()) == null) { - spans.put(abstractDataSerializable.getDataType(), new ArrayList()); - } - spans.get(abstractDataSerializable.getDataType()).add(abstractDataSerializable); - } - - return spans; - } - private File chooseDealBufferFile() { File file1 = null; File parentDir = new File(Config.Buffer.DATA_BUFFER_FILE_PARENT_DIR);