From 0dadf27be7a8948e520d48e9587769210bc0ca4e Mon Sep 17 00:00:00 2001 From: ascrutae Date: Fri, 4 Nov 2016 18:15:59 +0800 Subject: [PATCH] =?UTF-8?q?=E9=87=8D=E6=9E=84=E4=BC=AA=E4=BB=A3=E7=A0=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../storage/data/SpanDataConsumer.java | 22 ++++----- .../storage/data/file/DataFileWriter.java | 4 +- .../storage/data/index/IndexDBConnector.java | 30 ++++++++++++ ...rCache.java => IndexDBConnectorCache.java} | 12 ++--- .../data/index/IndexMetaCollections.java | 37 ++++++++++++++ .../storage/data/index/IndexMetaGroup.java | 49 +++++++++++++++++++ .../storage/data/index/IndexOperator.java | 41 +++++----------- .../data/index/IndexOperatorFactory.java | 20 -------- 8 files changed, 146 insertions(+), 69 deletions(-) create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnector.java rename skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/{IndexOperatorCache.java => IndexDBConnectorCache.java} (82%) create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaCollections.java create mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaGroup.java delete mode 100644 skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexOperatorFactory.java 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 8ee27865d..ffd3946ed 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 @@ -1,30 +1,28 @@ package com.a.eye.skywalking.storage.data; import com.a.eye.datacarrier.consumer.IConsumer; -import com.a.eye.skywalking.storage.block.index.BlockIndexEngine; 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.IndexMetaGroup; import com.a.eye.skywalking.storage.data.index.IndexOperator; -import com.a.eye.skywalking.storage.data.index.IndexOperatorCache; +import com.a.eye.skywalking.storage.data.index.IndexDBConnectorCache; +import java.util.Iterator; import java.util.List; -import java.util.Map; public class SpanDataConsumer implements IConsumer { - private IndexOperatorCache cache; - private DataFileWriter fileWriter; + private IndexDBConnectorCache cache; + private DataFileWriter fileWriter; @Override public void consume(List data) { - List indexMetaInfo = fileWriter.write(data); - Map> categorizedMetaInfo = - IndexMetaInfoCategory.categorizeByDataIndexTime(indexMetaInfo, BlockIndexEngine.newFinder()); + Iterator iterator = fileWriter.write(data).group().iterator(); - for (Map.Entry> indexEntry : categorizedMetaInfo.entrySet()) { - IndexOperator indexOperator = cache.get(indexEntry.getKey()); - indexOperator.update(indexEntry.getValue()); + while (iterator.hasNext()) { + IndexMetaGroup metaGroup = iterator.next(); + IndexOperator indexOperator = IndexOperator.newOperator(cache.get(metaGroup.getTimestamp())); + indexOperator.update(metaGroup.getMetaInfo()); } } 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 7f79804dc..0eaf3bd68 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.IndexMetaInfo; +import com.a.eye.skywalking.storage.data.index.IndexMetaCollections; import java.util.List; @@ -13,7 +13,7 @@ public class DataFileWriter { dataFile = DataFilesManager.createNewDataFile(); } - public List write(List spanData) { + public IndexMetaCollections write(List spanData) { if (dataFile.overLimitLength()) { dataFile = DataFilesManager.createNewDataFile(); 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 new file mode 100644 index 000000000..be1b56965 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnector.java @@ -0,0 +1,30 @@ +package com.a.eye.skywalking.storage.data.index; + +/** + * Created by xin on 2016/11/4. + */ +public class IndexDBConnector { + + private long timestamp; + + public IndexDBConnector(long timestamp) { + + } + + private void validate() { + + } + + private void createTable() { + + } + + private void createIndex() { + + } + + + public long getTimestamp() { + return timestamp; + } +} diff --git a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexOperatorCache.java b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnectorCache.java similarity index 82% rename from skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexOperatorCache.java rename to skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnectorCache.java index 477015fc3..a2474da4b 100644 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexOperatorCache.java +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexDBConnectorCache.java @@ -6,21 +6,21 @@ import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReadWriteLock; import java.util.concurrent.locks.ReentrantReadWriteLock; -public class IndexOperatorCache { +public class IndexDBConnectorCache { private static final int MAX_CACHE_SIZE = 5; - private LRUCache cachedOperators; + private LRUCache cachedOperators; - public IndexOperatorCache() { - cachedOperators = new LRUCache(MAX_CACHE_SIZE); + public IndexDBConnectorCache() { + cachedOperators = new LRUCache(MAX_CACHE_SIZE); } - public IndexOperator get(long timestamp) { + public IndexDBConnector get(long timestamp) { return cachedOperators.get(timestamp); } - public void updateCache(long timestamp, IndexOperator operator) { + public void updateCache(long timestamp, IndexDBConnector operator) { cachedOperators.put(timestamp, operator); } 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 new file mode 100644 index 000000000..4e9e16cc2 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaCollections.java @@ -0,0 +1,37 @@ +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.List; + +public class IndexMetaCollections { + + private List metaInfo; + private BlockFinder finder = BlockIndexEngine.newFinder(); + + public List 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; + } + + +} 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 new file mode 100644 index 000000000..e791a5498 --- /dev/null +++ b/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexMetaGroup.java @@ -0,0 +1,49 @@ +package com.a.eye.skywalking.storage.data.index; + +import java.util.ArrayList; +import java.util.List; + +/** + * Created by xin on 2016/11/4. + */ +public class IndexMetaGroup { + + private long timestamp; + + private List metaInfo; + + public IndexMetaGroup(long timestamp) { + this.timestamp = timestamp; + metaInfo = new ArrayList(); + } + + public long getTimestamp() { + return timestamp; + } + + public List getMetaInfo() { + return metaInfo; + } + + public void addIndexMetaInfo(IndexMetaInfo info) { + this.metaInfo.add(info); + } + + @Override + public boolean equals(Object o) { + if (this == o) + return true; + if (o == null || getClass() != o.getClass()) + return false; + + IndexMetaGroup that = (IndexMetaGroup) o; + + return timestamp == that.timestamp; + + } + + @Override + public int hashCode() { + return (int) (timestamp ^ (timestamp >>> 32)); + } +} 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 index e6659906e..6503a6b9a 100644 --- 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 @@ -1,18 +1,16 @@ package com.a.eye.skywalking.storage.data.index; -import java.sql.Connection; import java.util.ArrayList; import java.util.List; -import static com.a.eye.skywalking.storage.config.Config.DataIndex.TABLE_NAME; - public class IndexOperator { - private Connection connection; - private long timestamp; - - private IndexOperator(long timestamp) { + private IndexDBConnector connector; + private long timestamp; + private IndexOperator(IndexDBConnector connector) { + this.connector = connector; + timestamp = connector.getTimestamp(); } public List find(String taceId) { @@ -24,31 +22,16 @@ public class IndexOperator { } - private Connection getConnection() { - return connection; + private IndexDBConnector getConnector() { + return connector; } - public static class Builder { - private IndexOperator operator; - private IndexOperatorHelper indexOperatorHelper; - - private Builder(long timestamp) { - operator = new IndexOperator(timestamp); - indexOperatorHelper = new IndexOperatorHelper(operator.getConnection()); - } - - public static Builder newBuilder(long timestamp) { - return new Builder(timestamp); - } - - public IndexOperator build() { - if (indexOperatorHelper.validateIsReady(TABLE_NAME)) { - indexOperatorHelper.maintain(); - } - - return operator; - } + public static IndexOperator newOperator(long timestamp) { + return newOperator(new IndexDBConnector(timestamp)); } + public static IndexOperator newOperator(IndexDBConnector indexDBConnector) { + return new IndexOperator(indexDBConnector); + } } 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 deleted file mode 100644 index 11d17b986..000000000 --- a/skywalking-storage-center/skywalking-storage/src/main/java/com/a/eye/skywalking/storage/data/index/IndexOperatorFactory.java +++ /dev/null @@ -1,20 +0,0 @@ -package com.a.eye.skywalking.storage.data.index; - -/** - * Created by xin on 2016/11/4. - */ -public class IndexOperatorFactory { - - private static IndexOperatorCache operatorCache; - - public static IndexOperator getIndexOperator(long timestamp) { - IndexOperator operator = operatorCache.get(timestamp); - - if (operator == null) { - operator = IndexOperator.Builder.newBuilder(timestamp).build(); - operatorCache.updateCache(timestamp, operator); - } - - return operator; - } -}