From c06b1fd0a76033d12885c96a897dc8ff8576b0fe Mon Sep 17 00:00:00 2001 From: ascrutae Date: Tue, 15 Nov 2016 15:52:25 +0800 Subject: [PATCH] =?UTF-8?q?=E6=96=B0=E5=A2=9E=E9=85=8D=E7=BD=AE=E9=A1=B9?= =?UTF-8?q?=EF=BC=8C=E6=B7=BB=E5=8A=A0=E6=B5=8B=E8=AF=95=E7=94=A8=E4=BE=8B?= =?UTF-8?q?=EF=BC=8C=E4=BF=AE=E5=A4=8D=E5=9C=A8=E5=BF=AB=E9=80=9F=E5=88=9B?= =?UTF-8?q?=E5=BB=BADataFile=E6=97=B6=EF=BC=8C=E4=BC=9A=E5=88=9B=E5=BB=BA?= =?UTF-8?q?=E5=87=BA=E4=B8=A4=E4=B8=AA=E7=9B=B8=E5=90=8C=E5=90=8D=E5=AD=97?= =?UTF-8?q?=E7=9A=84DataFile?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../com/a/eye/skywalking/storage/Main.java | 2 +- .../storage/block/index/BlockFinder.java | 4 - .../eye/skywalking/storage/config/Config.java | 10 ++- .../storage/data/file/DataFile.java | 17 +++- .../storage/listener/StorageListener.java | 4 +- .../src/main/resources/config.properties | 9 +++ .../src/test/java/StorageClient.java | 78 +++++++++++++++++++ .../storage/data/file/DataFileWriterTest.java | 58 ++++++++++++++ 8 files changed, 170 insertions(+), 12 deletions(-) create mode 100644 skywalking-storage-center/skywalking-storage/src/test/java/StorageClient.java create mode 100644 skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/data/file/DataFileWriterTest.java diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/Main.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/Main.java index f658a2099..afe9e0783 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/Main.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/Main.java @@ -30,7 +30,7 @@ import static com.a.eye.skywalking.storage.config.Config.RegistryCenter.PATH_PRE public class Main { private static final ILog logger = LogManager.getLogger(Main.class); - private static final String SERVER_REPORTER_NAME = "Storage Server"; + private static final String SERVER_REPORTER_NAME = "DataConsumer Server"; static { LogManager.setLogResolver(new Log4j2Resolver()); diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/block/index/BlockFinder.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/block/index/BlockFinder.java index 90f1fe690..165d6e0b6 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/block/index/BlockFinder.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/block/index/BlockFinder.java @@ -24,10 +24,6 @@ public class BlockFinder { index = l2Cache.find(timestamp); } - if (logger.isDebugEnable()) { - logger.debug("Time stamp : {} is mapping with block Index {}.", timestamp, index); - } - return index; } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java index 3740b57ae..7c123d1ea 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/config/Config.java @@ -6,10 +6,16 @@ package com.a.eye.skywalking.storage.config; public class Config { public static class Server { public static int PORT = 34000; + } - public static int CHANNEL_SIZE = 10; - public static int BUFFER_SIZE = 1000; + public static class DataConsumer { + + public static int CHANNEL_SIZE = 10; + + public static int BUFFER_SIZE = 1000; + + public static int CONSUMER_SIZE = 5; } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFile.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFile.java index 9721663fe..327878477 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFile.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFile.java @@ -1,5 +1,8 @@ package com.a.eye.skywalking.storage.data.file; +import com.a.eye.datacarrier.common.AtomicRangeInteger; +import com.a.eye.skywalking.health.report.HealthCollector; +import com.a.eye.skywalking.health.report.HeathReading; import com.a.eye.skywalking.logging.api.ILog; import com.a.eye.skywalking.logging.api.LogManager; import com.a.eye.skywalking.storage.config.Config; @@ -13,6 +16,9 @@ import java.io.File; import java.io.FileInputStream; import java.io.FileOutputStream; import java.io.IOException; +import java.text.SimpleDateFormat; +import java.util.Date; +import java.util.concurrent.atomic.AtomicInteger; import static com.a.eye.skywalking.storage.util.PathResolver.getAbsolutePath; @@ -25,6 +31,7 @@ public class DataFile { private String fileName; private long currentOffset; private DataFileOperator operator; + private static final AtomicRangeInteger DATA_FILE_NAME_SUFFIX = new AtomicRangeInteger(1000, 9999); static { File dataFileDir = new File(getAbsolutePath(Config.DataFile.PATH)); @@ -34,7 +41,8 @@ public class DataFile { } public DataFile() { - this.fileName = System.currentTimeMillis() + ""; + this.fileName = new SimpleDateFormat("yyyy_MM_dd_HH_mm_ss_SS").format(new Date()) + "_" + DATA_FILE_NAME_SUFFIX + .getAndIncrement(); this.currentOffset = 0; operator = new DataFileOperator(); createFile(); @@ -61,6 +69,8 @@ public class DataFile { if (logger.isDebugEnable()) { logger.debug("Create an new data file[{}].", fileName); } + HealthCollector.getCurrentHeathReading("DataFile") + .updateData(HeathReading.INFO, "Create an new data " + "file."); } catch (IOException e) { logger.error("Failed to create data file.", e); throw new DataFileOperatorCreateFailedException("Failed to create data file", e); @@ -96,7 +106,7 @@ public class DataFile { } } - public void close(){ + public void close() { operator.close(); } @@ -107,7 +117,8 @@ public class DataFile { operator.getReader().read(data, 0, length); return data; } catch (IOException e) { - throw new SpanDataReadFailedException("Failed to read dataFile[" + fileName + "], offset: " + offset + " " + "lenght: " + length, e); + throw new SpanDataReadFailedException( + "Failed to read dataFile[" + fileName + "], offset: " + offset + " " + "lenght: " + length, e); } } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/StorageListener.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/StorageListener.java index bc24ff6c7..0b5e8e4ef 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/StorageListener.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/listener/StorageListener.java @@ -20,8 +20,8 @@ public class StorageListener implements SpanStorageListener { private DataCarrier spanDataDataCarrier; public StorageListener() { - spanDataDataCarrier = new DataCarrier<>(Config.Server.CHANNEL_SIZE, Config.Server.BUFFER_SIZE); - spanDataDataCarrier.consume(new SpanDataConsumer(), 5, true); + spanDataDataCarrier = new DataCarrier<>(Config.DataConsumer.CHANNEL_SIZE, Config.DataConsumer.BUFFER_SIZE); + spanDataDataCarrier.consume(new SpanDataConsumer(), Config.DataConsumer.CONSUMER_SIZE, true); } @Override diff --git a/skywalking-storage-center/skywalking-storage/src/main/resources/config.properties b/skywalking-storage-center/skywalking-storage/src/main/resources/config.properties index 12ef02c05..6a88961bc 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/resources/config.properties +++ b/skywalking-storage-center/skywalking-storage/src/main/resources/config.properties @@ -1,6 +1,15 @@ # the port which storage server listening server.port=34000 # +# the size of channel which storage span data +#dataconsumer.channel_size = 10 +# +# the buffer size for each channel +#dataconsumer.buffer_size = 1000 +# +# the size of data consumer +#dataconsumer.consumer_size = 5 +# # the path that storage block index #blockindex.path=/block-index # diff --git a/skywalking-storage-center/skywalking-storage/src/test/java/StorageClient.java b/skywalking-storage-center/skywalking-storage/src/test/java/StorageClient.java new file mode 100644 index 000000000..e69a4b03a --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/test/java/StorageClient.java @@ -0,0 +1,78 @@ +import com.a.eye.skywalking.network.dependencies.io.grpc.ManagedChannel; +import com.a.eye.skywalking.network.dependencies.io.grpc.ManagedChannelBuilder; +import com.a.eye.skywalking.network.dependencies.io.grpc.stub.StreamObserver; +import com.a.eye.skywalking.network.grpc.AckSpan; +import com.a.eye.skywalking.network.grpc.RequestSpan; +import com.a.eye.skywalking.network.grpc.SendResult; +import com.a.eye.skywalking.network.grpc.SpanStorageServiceGrpc; + +import static com.a.eye.skywalking.network.grpc.SpanStorageServiceGrpc.newStub; + +public class StorageClient { + private static ManagedChannel channel = + ManagedChannelBuilder.forAddress("127.0.0.1", 34000).usePlaintext(true).build(); + + private static SpanStorageServiceGrpc.SpanStorageServiceStub spanStorageServiceStub = newStub(channel); + + private static StreamObserver ackSpanStreamObserver = + spanStorageServiceStub.storageACKSpan(new StreamObserver() { + @Override + public void onNext(SendResult sendResult) { + System.out.println(sendResult.getResult()); + } + + @Override + public void onError(Throwable throwable) { + throwable.printStackTrace(); + } + + @Override + public void onCompleted() { + System.out.println("Success!!"); + } + }); + + + private static StreamObserver requestSpanStreamObserver = + spanStorageServiceStub.storageRequestSpan(new StreamObserver() { + @Override + public void onNext(SendResult sendResult) { + System.out.println(sendResult.getResult()); + } + + @Override + public void onError(Throwable throwable) { + throwable.printStackTrace(); + } + + @Override + public void onCompleted() { + System.out.println("Success!!"); + } + }); + + + public static void main(String[] args) throws InterruptedException { + RequestSpan requestSpan = + RequestSpan.newBuilder().setSpanType(1).setAddress("127.0.0.1").setApplicationId("1").setCallType("1") + .setLevelId(0).setProcessNo("19287").setStartDate(System.currentTimeMillis()) + .setTraceId("1.0Final.1478661327960.8504828.2277.53.3").setUserId("1") + .setViewPointId("http://localhost:8080/wwww/test/helloWorld").build(); + AckSpan ackSpan = + AckSpan.newBuilder().setLevelId(0).setCost(10).setTraceId("1.0Final.1478661327960.8504828.2277.53.3") + .setStatusCode(0).setViewpointId("http://localhost:8080/wwww/test/helloWorld").build(); + + for (int i = 0; i < 100000; i++) { + requestSpanStreamObserver.onNext(requestSpan); + ackSpanStreamObserver.onNext(ackSpan); + Thread.sleep(100); + } + + + ackSpanStreamObserver.onCompleted(); + requestSpanStreamObserver.onCompleted(); + + Thread.sleep(10000); + + } +} diff --git a/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/data/file/DataFileWriterTest.java b/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/data/file/DataFileWriterTest.java new file mode 100644 index 000000000..0bccd2011 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/test/java/com/a/eye/skywalking/storage/data/file/DataFileWriterTest.java @@ -0,0 +1,58 @@ +package com.a.eye.skywalking.storage.data.file; + +import com.a.eye.skywalking.network.grpc.RequestSpan; +import com.a.eye.skywalking.storage.config.Config; +import com.a.eye.skywalking.storage.data.spandata.RequestSpanData; +import com.a.eye.skywalking.storage.data.spandata.SpanData; +import com.a.eye.skywalking.storage.util.PathResolver; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +import java.io.File; +import java.io.IOException; +import java.nio.file.Files; +import java.util.ArrayList; +import java.util.List; + +import static org.junit.Assert.assertEquals; + +/** + * Created by xin on 2016/11/15. + */ +public class DataFileWriterTest { + + private DataFileWriter writer; + + @Before + public void setUp() { + Config.DataFile.PATH = "/tmp"; + Config.DataFile.SIZE = 10; + writer = new DataFileWriter(); + } + + @Test + public void testConvertFile() throws Exception { + List spanData = new ArrayList<>(); + spanData.add(new RequestSpanData( + RequestSpan.newBuilder().setTraceId("test-traceId").setStartDate(System.currentTimeMillis()) + .setProcessNo("7777").setLevelId(10).setParentLevel("0.0.0").setAddress("127.0.0.1").build())); + writer.write(spanData); + + writer.write(spanData); + File dir = new File(PathResolver.getAbsolutePath(Config.DataFile.PATH)); + assertEquals(2, dir.listFiles().length); + } + + + @After + public void tearUp() throws IOException { + File dir = new File(PathResolver.getAbsolutePath(Config.DataFile.PATH)); + for (File file : dir.listFiles()) { + file.delete(); + } + + dir.delete(); + } + +}