diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataConsumer.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataConsumer.java index c5034e2dd..57a555e6e 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataConsumer.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataConsumer.java @@ -2,10 +2,7 @@ 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.IndexDBConnector; -import com.a.eye.skywalking.storage.data.index.IndexMetaGroup; -import com.a.eye.skywalking.storage.data.index.IndexOperator; -import com.a.eye.skywalking.storage.data.index.IndexDBConnectorCache; +import com.a.eye.skywalking.storage.data.index.*; import java.util.Iterator; import java.util.List; @@ -15,7 +12,7 @@ public class SpanDataConsumer implements IConsumer { private IndexDBConnectorCache cache; private DataFileWriter fileWriter; - public SpanDataConsumer(){ + public SpanDataConsumer() { cache = new IndexDBConnectorCache(); fileWriter = new DataFileWriter(); } @@ -23,18 +20,24 @@ public class SpanDataConsumer implements IConsumer { @Override public void consume(List data) { - Iterator iterator = fileWriter.write(data).group(); + Iterator> iterator = + IndexMetaCollections.group(fileWriter.write(data), new GroupKeyBuilder() { + @Override + public Long buildKey(IndexMetaInfo metaInfo) { + return metaInfo.getStartTime(); + } + }).iterator(); while (iterator.hasNext()) { - IndexMetaGroup metaGroup = iterator.next(); + IndexMetaGroup metaGroup = iterator.next(); IndexOperator indexOperator = IndexOperator.newOperator(getDBConnector(metaGroup)); indexOperator.batchUpdate(metaGroup); } } - private IndexDBConnector getDBConnector(IndexMetaGroup metaGroup){ - return cache.get(metaGroup.getTimestamp()); + private IndexDBConnector getDBConnector(IndexMetaGroup metaGroup) { + return cache.get(metaGroup.getKey()); } @Override diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataFinder.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataFinder.java new file mode 100644 index 000000000..9dfd66d3d --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/SpanDataFinder.java @@ -0,0 +1,42 @@ +package com.a.eye.skywalking.storage.data; + +import com.a.eye.skywalking.storage.block.index.BlockIndexEngine; +import com.a.eye.skywalking.storage.data.file.DataFileReader; +import com.a.eye.skywalking.storage.data.file.DataFileWriter; +import com.a.eye.skywalking.storage.data.index.*; + +import java.util.ArrayList; +import java.util.Iterator; +import java.util.List; + +/** + * Created by xin on 2016/11/6. + */ +public class SpanDataFinder { + + public static List find(String traceId) { + long blockIndex = BlockIndexEngine.newFinder().find(fetchStartTimeFromTraceId(traceId)); + IndexDBConnector indexDBConnector = new IndexDBConnector(blockIndex); + IndexMetaCollection indexMetaCollection = indexDBConnector.queryByTraceId(traceId); + + Iterator> iterator = + IndexMetaCollections.group(indexMetaCollection, new GroupKeyBuilder() { + @Override + public String buildKey(IndexMetaInfo metaInfo) { + return metaInfo.getFileName(); + } + }).iterator(); + + List result = new ArrayList(); + while (iterator.hasNext()) { + IndexMetaGroup group = iterator.next(); + result.addAll(new DataFileReader(group.getKey()).read(group.getMetaInfo())); + } + + return result; + } + + private static long fetchStartTimeFromTraceId(String traceId) { + return -1; + } +} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/exception/SpanDataReadFailedException.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/exception/SpanDataReadFailedException.java new file mode 100644 index 000000000..6de32e2f3 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/exception/SpanDataReadFailedException.java @@ -0,0 +1,10 @@ +package com.a.eye.skywalking.storage.data.exception; + +/** + * Created by xin on 2016/11/6. + */ +public class SpanDataReadFailedException extends RuntimeException { + public SpanDataReadFailedException(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 index 1ce3f94c1..d057edadc 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 @@ -4,6 +4,7 @@ import com.a.eye.skywalking.storage.config.Config; import com.a.eye.skywalking.storage.data.SpanData; import com.a.eye.skywalking.storage.data.exception.DataFileOperatorCreateFailedException; import com.a.eye.skywalking.storage.data.exception.SpanDataPersistenceFailedException; +import com.a.eye.skywalking.storage.data.exception.SpanDataReadFailedException; import com.a.eye.skywalking.storage.data.index.IndexMetaInfo; import java.io.File; @@ -26,9 +27,8 @@ public class DataFile { operator = new DataFileOperator(); } - public DataFile(String fileName, long offset) { + public DataFile(String fileName) { this.fileName = fileName; - this.currentOffset = offset; operator = new DataFileOperator(); } @@ -62,6 +62,18 @@ public class DataFile { } } + public byte[] read(long offset, int length) { + byte[] data = new byte[length]; + try { + operator.getReader().getChannel().position(offset); + operator.getReader().read(data, 0, length); + return data; + } catch (IOException e) { + throw new SpanDataReadFailedException( + "Failed to read dataFile[" + fileName + "], offset: " + offset + " " + "lenght: " + length, e); + } + } + class DataFileOperator { private FileOutputStream writer; private FileInputStream reader; 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 new file mode 100644 index 000000000..c0b4ba40a --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/file/DataFileReader.java @@ -0,0 +1,28 @@ +package com.a.eye.skywalking.storage.data.file; + +import com.a.eye.skywalking.storage.data.index.IndexMetaInfo; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; + +/** + * Created by xin on 2016/11/6. + */ +public class DataFileReader { + private DataFile dataFile; + + public DataFileReader(String fileName) { + dataFile = new DataFile(fileName); + } + + public List read(List metaInfo) { + List metaData = new ArrayList(); + + for (IndexMetaInfo indexMetaInfo : metaInfo){ + metaData.add(dataFile.read(indexMetaInfo.getOffset(), indexMetaInfo.getLength())); + } + + return metaData; + } +} 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 4e6c34281..a6878bc83 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,7 @@ package com.a.eye.skywalking.storage.data.file; import com.a.eye.skywalking.storage.data.SpanData; -import com.a.eye.skywalking.storage.data.index.IndexMetaCollections; +import com.a.eye.skywalking.storage.data.index.IndexMetaCollection; import java.util.List; @@ -13,12 +13,12 @@ public class DataFileWriter { dataFile = DataFilesManager.createNewDataFile(); } - public IndexMetaCollections write(List spanData) { + public IndexMetaCollection write(List spanData) { if (dataFile.overLimitLength()) { dataFile = DataFilesManager.createNewDataFile(); } - IndexMetaCollections collections = new IndexMetaCollections(); + IndexMetaCollection collections = new IndexMetaCollection(); for (SpanData data : spanData) { collections.add(dataFile.write(data)); } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/GroupKeyBuilder.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/GroupKeyBuilder.java new file mode 100644 index 000000000..086577e3f --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/GroupKeyBuilder.java @@ -0,0 +1,10 @@ +package com.a.eye.skywalking.storage.data.index; + +import com.a.eye.skywalking.storage.data.index.IndexMetaInfo; + +/** + * Created by xin on 2016/11/6. + */ +public interface GroupKeyBuilder { + T buildKey(IndexMetaInfo metaInfo); +} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnector.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnector.java index 89c8e0a9a..c021189ec 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnector.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnector.java @@ -86,7 +86,7 @@ public class IndexDBConnector { return timestamp; } - public void batchUpdate(IndexMetaGroup metaGroup) throws SQLException { + public void batchUpdate(IndexMetaGroup metaGroup) throws SQLException { int currentIndex = 0; PreparedStatement ps = connection.prepareStatement(INSERT_INDEX); for (IndexMetaInfo metaInfo : metaGroup.getMetaInfo()) { @@ -117,6 +117,10 @@ public class IndexDBConnector { return indexSize; } + public IndexMetaCollection queryByTraceId(String traceId) { + return null; + } + class ConnectURLGenerator { private String basePath; diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaCollection.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaCollection.java new file mode 100644 index 000000000..33778effe --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaCollection.java @@ -0,0 +1,51 @@ +package com.a.eye.skywalking.storage.data.index; + + +import com.a.eye.skywalking.storage.block.index.BlockFinder; +import com.a.eye.skywalking.storage.block.index.BlockIndexEngine; + +import java.util.ArrayList; +import java.util.Iterator; +import java.util.List; + +public class IndexMetaCollection implements Iterable { + + private List metaInfo; + private BlockFinder finder; + + public IndexMetaCollection() { + metaInfo = new ArrayList<>(); + finder = BlockIndexEngine.newFinder(); + } + + public Iterator group() { + List indexMetaGroups = new ArrayList(); + for (IndexMetaInfo info : metaInfo) { + long timestamp = finder.find(info.getStartTime()); + + int index = indexMetaGroups.indexOf(new IndexMetaGroup(timestamp)); + IndexMetaGroup metaGroup; + + if (index == -1) { + metaGroup = new IndexMetaGroup(timestamp); + indexMetaGroups.add(metaGroup); + } else { + metaGroup = indexMetaGroups.get(index); + } + + metaGroup.addIndexMetaInfo(info); + } + + return indexMetaGroups.iterator(); + } + + + public void add(IndexMetaInfo info) { + metaInfo.add(info); + } + + @Override + public Iterator iterator() { + return metaInfo.iterator(); + } +} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaCollections.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaCollections.java index 40c92f2dd..ad229ed32 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaCollections.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaCollections.java @@ -1,46 +1,33 @@ package com.a.eye.skywalking.storage.data.index; - -import com.a.eye.skywalking.storage.block.index.BlockFinder; -import com.a.eye.skywalking.storage.block.index.BlockIndexEngine; - import java.util.ArrayList; -import java.util.Iterator; import java.util.List; +/** + * Created by xin on 2016/11/6. + */ public class IndexMetaCollections { - private List metaInfo; - private BlockFinder finder; + public static List> group(IndexMetaCollection indexMetaCollection, + GroupKeyBuilder builder) { + List> indexMetaGroups = new ArrayList>(); - public IndexMetaCollections() { - metaInfo = new ArrayList<>(); - finder = BlockIndexEngine.newFinder(); - } + for (IndexMetaInfo metaInfo : indexMetaCollection) { + T key = builder.buildKey(metaInfo); - public Iterator group() { - List indexMetaGroups = new ArrayList(); - for (IndexMetaInfo info : metaInfo) { - long timestamp = finder.find(info.getStartTime()); - - int index = indexMetaGroups.indexOf(new IndexMetaGroup(timestamp)); + int index = indexMetaGroups.indexOf(new IndexMetaGroup(key)); IndexMetaGroup metaGroup; if (index == -1) { - metaGroup = new IndexMetaGroup(timestamp); + metaGroup = new IndexMetaGroup(key); indexMetaGroups.add(metaGroup); } else { metaGroup = indexMetaGroups.get(index); } - metaGroup.addIndexMetaInfo(info); + metaGroup.addIndexMetaInfo(metaInfo); } - return indexMetaGroups.iterator(); - } - - - public void add(IndexMetaInfo info) { - metaInfo.add(info); + return indexMetaGroups; } } diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaGroup.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaGroup.java index 1e4c90b2b..95f4fa384 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaGroup.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaGroup.java @@ -1,24 +1,25 @@ package com.a.eye.skywalking.storage.data.index; import java.util.ArrayList; +import java.util.Iterator; import java.util.List; /** * Created by xin on 2016/11/4. */ -public class IndexMetaGroup { +public class IndexMetaGroup{ - private long timestamp; + private V key; private List metaInfo; - public IndexMetaGroup(long timestamp) { - this.timestamp = timestamp; + public IndexMetaGroup(V key) { + this.key = key; metaInfo = new ArrayList(); } - public long getTimestamp() { - return timestamp; + public V getKey() { + return key; } public List getMetaInfo() { @@ -36,18 +37,19 @@ public class IndexMetaGroup { if (o == null || getClass() != o.getClass()) return false; - IndexMetaGroup that = (IndexMetaGroup) o; + IndexMetaGroup that = (IndexMetaGroup) o; - return timestamp == that.timestamp; + return key.equals(that.key); } @Override public int hashCode() { - return (int) (timestamp ^ (timestamp >>> 32)); + return key.hashCode(); } public int size() { return metaInfo.size(); } + }