diff --git a/skywalking-storage-center/skywalking-storage/pom.xml b/skywalking-storage-center/skywalking-storage/pom.xml
index b51d698e1..98c6883ab 100644
--- a/skywalking-storage-center/skywalking-storage/pom.xml
+++ b/skywalking-storage-center/skywalking-storage/pom.xml
@@ -67,6 +67,11 @@
2.3.4
+
+ com.a.eye
+ data-carrier
+ 1.0
+
@@ -159,5 +164,14 @@
-
+
+
+
+ false
+
+ bintray-wu-sheng-DataCarrier
+ bintray
+ http://dl.bintray.com/wu-sheng/DataCarrier
+
+
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 186b7fa68..6e67945ee 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
@@ -10,4 +10,9 @@ public class Config {
public static String DATA_FILE_INDEX_FILE_NAME = "data_file.index";
}
+
+
+ public static class DataFile {
+ public static String BASE_PATH = "";
+ }
}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/DataOperator.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/DataOperator.java
deleted file mode 100644
index 8cd4d6566..000000000
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/DataOperator.java
+++ /dev/null
@@ -1,8 +0,0 @@
-package com.a.eye.skywalking.storage.data;
-
-/**
- * Created by xin on 2016/11/1.
- */
-public class DataOperator {
-
-}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanData.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanData.java
new file mode 100644
index 000000000..7461daf47
--- /dev/null
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanData.java
@@ -0,0 +1,10 @@
+package com.a.eye.skywalking.storage.data;
+
+/**
+ * Created by xin on 2016/11/4.
+ */
+public interface SpanData {
+ long getTimestamp();
+
+ byte[] convertToByte();
+}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataBuilder.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataBuilder.java
new file mode 100644
index 000000000..9ea41d14d
--- /dev/null
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataBuilder.java
@@ -0,0 +1,7 @@
+package com.a.eye.skywalking.storage.data;
+
+/**
+ * Created by xin on 2016/11/4.
+ */
+public class SpanDataBuilder {
+}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataWriterConsumer.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataWriterConsumer.java
new file mode 100644
index 000000000..45466af03
--- /dev/null
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataWriterConsumer.java
@@ -0,0 +1,28 @@
+package com.a.eye.skywalking.storage.data;
+
+import com.a.eye.datacarrier.consumer.IConsumer;
+import com.a.eye.skywalking.storage.data.file.DataFileWriter;
+import com.a.eye.skywalking.storage.data.index.IndexMetaInfo;
+import com.a.eye.skywalking.storage.data.index.IndexOperator;
+import com.a.eye.skywalking.storage.data.index.IndexOperatorFactory;
+
+import java.util.List;
+
+public class SpanDataWriterConsumer implements IConsumer {
+
+ private DataFileWriter writer;
+
+ @Override
+ public void consume(List list) {
+ for (SpanData data : list) {
+ IndexOperator operator = IndexOperatorFactory.get(data.getTimestamp());
+ IndexMetaInfo metaInfo = writer.write(data.convertToByte());
+ operator.update(metaInfo);
+ }
+ }
+
+ @Override
+ public void onError(List list, Throwable throwable) {
+
+ }
+}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/exception/DataFileNotFoundException.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/exception/DataFileNotFoundException.java
new file mode 100644
index 000000000..02cd791c1
--- /dev/null
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/exception/DataFileNotFoundException.java
@@ -0,0 +1,10 @@
+package com.a.eye.skywalking.storage.data.exception;
+
+/**
+ * Created by xin on 2016/11/4.
+ */
+public class DataFileNotFoundException extends Exception {
+ public DataFileNotFoundException(String message, Exception e) {
+ super(message, e);
+ }
+}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/exception/FileReaderCreateFailedException.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/exception/FileReaderCreateFailedException.java
new file mode 100644
index 000000000..6511ac029
--- /dev/null
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/exception/FileReaderCreateFailedException.java
@@ -0,0 +1,10 @@
+package com.a.eye.skywalking.storage.data.exception;
+
+/**
+ * Created by xin on 2016/11/4.
+ */
+public class FileReaderCreateFailedException extends Throwable {
+ public FileReaderCreateFailedException(String message, Exception e) {
+ super(message, e);
+ }
+}
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
new file mode 100644
index 000000000..6a58f99b9
--- /dev/null
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFile.java
@@ -0,0 +1,27 @@
+package com.a.eye.skywalking.storage.data.file;
+
+public class DataFile {
+
+ private static final long MAX_LENGTH = 3 * 1024 * 1024 * 1024;
+
+ private String name;
+ private long currentLength;
+
+ public boolean overLimitLength() {
+ return currentLength > MAX_LENGTH;
+ }
+
+
+ public long writeAndFlush(byte[] data) {
+
+ return 0;
+ }
+
+ public String getName() {
+ return name;
+ }
+
+ public void close() {
+
+ }
+}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileOperator.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileOperator.java
deleted file mode 100644
index d8573982f..000000000
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileOperator.java
+++ /dev/null
@@ -1,16 +0,0 @@
-package com.a.eye.skywalking.storage.data.file;
-
-/**
- * Created by xin on 2016/10/31.
- */
-public class DataFileOperator {
-
- public int write(){
-
- return 0;
- }
-
- public Object read(Object dataIndex) {
- return null;
- }
-}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileOperatorFactory.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileOperatorFactory.java
new file mode 100644
index 000000000..5e5e2b6dc
--- /dev/null
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileOperatorFactory.java
@@ -0,0 +1,20 @@
+package com.a.eye.skywalking.storage.data.file;
+
+import com.a.eye.skywalking.storage.data.exception.DataFileNotFoundException;
+import com.a.eye.skywalking.storage.data.exception.FileReaderCreateFailedException;
+import com.a.eye.skywalking.storage.data.index.IndexMetaInfo;
+
+public class DataFileOperatorFactory {
+
+ public static DataFileReader newReader(IndexMetaInfo indexMetaInfo) throws FileReaderCreateFailedException {
+ try {
+ return new DataFileReader(indexMetaInfo);
+ } catch (DataFileNotFoundException e) {
+ throw new FileReaderCreateFailedException("Cannot create DataFileReader.", e);
+ }
+ }
+
+ public static DataFileWriter newWriter() {
+ return new DataFileWriter();
+ }
+}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileReader.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileReader.java
index 864e8ca5c..53d0ebd43 100644
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileReader.java
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileReader.java
@@ -1,7 +1,45 @@
package com.a.eye.skywalking.storage.data.file;
-/**
- * Created by xin on 2016/10/31.
- */
+import com.a.eye.skywalking.storage.config.Config;
+import com.a.eye.skywalking.storage.data.exception.DataFileNotFoundException;
+import com.a.eye.skywalking.storage.data.index.IndexMetaInfo;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.io.File;
+import java.io.FileInputStream;
+import java.io.FileNotFoundException;
+import java.io.IOException;
+
public class DataFileReader {
+
+ private static Logger logger = LogManager.getLogger(DataFileReader.class);
+
+ private long offset;
+ private int length;
+ private String fileName;
+ private FileInputStream reader;
+
+ public DataFileReader(IndexMetaInfo indexMetaInfo) throws DataFileNotFoundException {
+ try {
+ reader = new FileInputStream(new File(Config.DataFile.BASE_PATH, indexMetaInfo.getFileName()));
+ this.offset = indexMetaInfo.getOffset();
+ this.length = indexMetaInfo.getLength();
+ this.fileName = indexMetaInfo.getFileName();
+ } catch (FileNotFoundException e) {
+ throw new DataFileNotFoundException(indexMetaInfo.getFileName() + " not found.", e);
+ }
+ }
+
+ public byte[] read() {
+ try {
+ reader.getChannel().position(this.offset);
+ byte[] dataByte = new byte[length];
+ reader.read(dataByte, 0, length);
+ return dataByte;
+ } catch (IOException e) {
+ logger.error("Failed to read file:{} position:{} length:{}", fileName, offset, length);
+ return null;
+ }
+ }
}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java
index 27751ebf4..f627390a6 100644
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileWriter.java
@@ -1,7 +1,27 @@
package com.a.eye.skywalking.storage.data.file;
-/**
- * Created by xin on 2016/10/31.
- */
+import com.a.eye.skywalking.storage.data.index.IndexMetaInfo;
+
public class DataFileWriter {
+
+ private DataFile dataFile;
+
+ public DataFileWriter() {
+
+ }
+
+ public IndexMetaInfo write(byte[] data) {
+ if (dataFile.overLimitLength()) {
+ convertDataFile();
+ }
+
+ long offset = dataFile.writeAndFlush(data);
+ return new IndexMetaInfo(dataFile.getName(), offset, data.length);
+ }
+
+ private void convertDataFile() {
+ dataFile.close();
+ dataFile = new DataFile();
+ }
+
}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/DataIndexFile.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/DataIndexFile.java
deleted file mode 100644
index 86c92cb84..000000000
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/DataIndexFile.java
+++ /dev/null
@@ -1,14 +0,0 @@
-package com.a.eye.skywalking.storage.data.index;
-
-/**
- * Created by xin on 2016/10/31.
- */
-public class DataIndexFile {
-
- public void read(String traceId) {
- }
-
- public void write() {
-
- }
-}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/DataIndexFileOperator.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/DataIndexFileOperator.java
deleted file mode 100644
index 6d8bc1946..000000000
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/DataIndexFileOperator.java
+++ /dev/null
@@ -1,16 +0,0 @@
-package com.a.eye.skywalking.storage.data.index;
-
-/**
- * Created by xin on 2016/10/31.
- */
-public class DataIndexFileOperator {
-
- public Object read(String traceId) {
- return null;
- }
-
- public void write(Object o, int offset) {
-
- }
-
-}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/DataIndexFileReader.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/DataIndexFileReader.java
deleted file mode 100644
index c06583c63..000000000
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/DataIndexFileReader.java
+++ /dev/null
@@ -1,7 +0,0 @@
-package com.a.eye.skywalking.storage.data.index;
-
-/**
- * Created by xin on 2016/10/31.
- */
-public class DataIndexFileReader {
-}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/DataIndexFileWriter.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/DataIndexFileWriter.java
deleted file mode 100644
index 320d64654..000000000
--- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/DataIndexFileWriter.java
+++ /dev/null
@@ -1,7 +0,0 @@
-package com.a.eye.skywalking.storage.data.index;
-
-/**
- * Created by xin on 2016/10/31.
- */
-public class DataIndexFileWriter {
-}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaInfo.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaInfo.java
new file mode 100644
index 000000000..27abeb362
--- /dev/null
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaInfo.java
@@ -0,0 +1,32 @@
+package com.a.eye.skywalking.storage.data.index;
+
+/**
+ * Created by xin on 2016/11/3.
+ */
+public class IndexMetaInfo {
+
+ private String fileName;
+
+ private long offset;
+
+ private int length;
+
+ public IndexMetaInfo(String name, long offset, int length) {
+ this.fileName = name;
+ this.offset = offset;
+ this.length = length;
+ }
+
+
+ public String getFileName() {
+ return fileName;
+ }
+
+ public long getOffset() {
+ return offset;
+ }
+
+ public int getLength() {
+ return length;
+ }
+}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexOperator.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexOperator.java
new file mode 100644
index 000000000..252b95895
--- /dev/null
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexOperator.java
@@ -0,0 +1,17 @@
+package com.a.eye.skywalking.storage.data.index;
+
+import java.util.List;
+
+/**
+ * Created by xin on 2016/11/3.
+ */
+public class IndexOperator {
+
+ public List find(String traceId) {
+ return null;
+ }
+
+ public void update(IndexMetaInfo meta) {
+
+ }
+}
diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexOperatorFactory.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexOperatorFactory.java
new file mode 100644
index 000000000..21b9da753
--- /dev/null
+++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexOperatorFactory.java
@@ -0,0 +1,8 @@
+package com.a.eye.skywalking.storage.data.index;
+
+public class IndexOperatorFactory {
+
+ public static IndexOperator get(long timestamp) {
+ return null;
+ }
+}