This is an automated email from the ASF dual-hosted git repository. jiangtian pushed a commit to branch dev_TTL in repository https://gitbox.apache.org/repos/asf/incubator-iotdb.git
commit 80e917a83325c0f3b3c652dfb67d655ad1206bb5 Author: jt <[email protected]> AuthorDate: Thu Sep 12 17:22:06 2019 +0800 add basic supports for TTL --- .../apache/iotdb/jdbc/IoTDBDatabaseMetadata.java | 2 +- .../apache/iotdb/jdbc/IoTDBMetadataResultSet.java | 26 +++---- .../org/apache/iotdb/db/engine/StorageEngine.java | 23 +++--- .../iotdb/db/engine/memtable/AbstractMemTable.java | 8 +-- .../apache/iotdb/db/engine/memtable/IMemTable.java | 2 +- .../db/engine/querycontext/QueryDataSource.java | 32 +++++++++ .../engine/storagegroup/StorageGroupProcessor.java | 83 ++++++++++++++++------ .../db/engine/storagegroup/TsFileProcessor.java | 11 ++- .../java/org/apache/iotdb/db/metadata/MGraph.java | 14 ++-- .../org/apache/iotdb/db/metadata/MManager.java | 66 +++++++++-------- .../java/org/apache/iotdb/db/metadata/MNode.java | 30 ++++++++ .../java/org/apache/iotdb/db/metadata/MTree.java | 26 ++++++- .../groupby/GroupByWithoutValueFilterDataSet.java | 1 + .../db/query/executor/AggregateEngineExecutor.java | 4 ++ .../SeriesReaderWithoutValueFilter.java | 9 +-- .../org/apache/iotdb/db/service/TSServiceImpl.java | 20 ++---- .../iotdb/db/metadata/MManagerAdvancedTest.java | 2 +- .../iotdb/db/metadata/MManagerBasicTest.java | 36 +++++----- .../java/org/apache/iotdb/rpc/TSStatusType.java | 1 + service-rpc/src/main/thrift/rpc.thrift | 2 +- 20 files changed, 259 insertions(+), 139 deletions(-) diff --git a/jdbc/src/main/java/org/apache/iotdb/jdbc/IoTDBDatabaseMetadata.java b/jdbc/src/main/java/org/apache/iotdb/jdbc/IoTDBDatabaseMetadata.java index 1344c2d..50f7cfd 100644 --- a/jdbc/src/main/java/org/apache/iotdb/jdbc/IoTDBDatabaseMetadata.java +++ b/jdbc/src/main/java/org/apache/iotdb/jdbc/IoTDBDatabaseMetadata.java @@ -112,7 +112,7 @@ public class IoTDBDatabaseMetadata implements DatabaseMetaData { } catch (IoTDBRPCException e) { throw new IoTDBSQLException(e.getMessage()); } - Set<String> showStorageGroup = resp.getShowStorageGroups(); + List<String> showStorageGroup = resp.getShowStorageGroups(); return new IoTDBMetadataResultSet(showStorageGroup, IoTDBMetadataResultSet.MetadataType.STORAGE_GROUP); } catch (TException e) { throw new TException("Conncetion error when fetching storage group metadata", e); diff --git a/jdbc/src/main/java/org/apache/iotdb/jdbc/IoTDBMetadataResultSet.java b/jdbc/src/main/java/org/apache/iotdb/jdbc/IoTDBMetadataResultSet.java index c52d327..e0170a6 100644 --- a/jdbc/src/main/java/org/apache/iotdb/jdbc/IoTDBMetadataResultSet.java +++ b/jdbc/src/main/java/org/apache/iotdb/jdbc/IoTDBMetadataResultSet.java @@ -25,16 +25,16 @@ import java.util.*; public class IoTDBMetadataResultSet extends IoTDBQueryResultSet { - public static final String GET_STRING_COLUMN = "COLUMN"; - public static final String GET_STRING_STORAGE_GROUP = "STORAGE_GROUP"; - public static final String GET_STRING_TIMESERIES_NUM = "TIMESERIES_NUM"; - public static final String GET_STRING_NODES_NUM = "NODE_NUM"; - public static final String GET_STRING_NODE_PATH = "NODE_PATH"; - public static final String GET_STRING_NODE_TIMESERIES_NUM = "NODE_TIMESERIES_NUM"; - public static final String GET_STRING_TIMESERIES_NAME = "Timeseries"; - public static final String GET_STRING_TIMESERIES_STORAGE_GROUP = "Storage Group"; + private static final String GET_STRING_COLUMN = "COLUMN"; + private static final String GET_STRING_STORAGE_GROUP = "STORAGE_GROUP"; + private static final String GET_STRING_TIMESERIES_NUM = "TIMESERIES_NUM"; + private static final String GET_STRING_NODES_NUM = "NODE_NUM"; + private static final String GET_STRING_NODE_PATH = "NODE_PATH"; + private static final String GET_STRING_NODE_TIMESERIES_NUM = "NODE_TIMESERIES_NUM"; + private static final String GET_STRING_TIMESERIES_NAME = "Timeseries"; + private static final String GET_STRING_TIMESERIES_STORAGE_GROUP = "Storage Group"; public static final String GET_STRING_TIMESERIES_DATATYPE = "DataType"; - public static final String GET_STRING_TIMESERIES_ENCODING = "Encoding"; + private static final String GET_STRING_TIMESERIES_ENCODING = "Encoding"; private Iterator<?> columnItr; private MetadataType type; private String currentColumn; @@ -56,7 +56,7 @@ public class IoTDBMetadataResultSet extends IoTDBQueryResultSet { /** * Constructor used for the result of DatabaseMetadata.getColumns() */ - public IoTDBMetadataResultSet(Object object, MetadataType type) throws SQLException { + IoTDBMetadataResultSet(Object object, MetadataType type) throws SQLException { this.type = type; switch (type) { case COLUMN: @@ -66,7 +66,7 @@ public class IoTDBMetadataResultSet extends IoTDBQueryResultSet { columnItr = columns.iterator(); break; case STORAGE_GROUP: - Set<String> storageGroupSet = (Set<String>) object; + List<String> storageGroupSet = (List<String>) object; colCount = 1; showLabels = new String[]{"Storage Group"}; columnItr = storageGroupSet.iterator(); @@ -79,7 +79,7 @@ public class IoTDBMetadataResultSet extends IoTDBQueryResultSet { columnItr = showTimeseriesList.iterator(); break; case COUNT_TIMESERIES: - String tsNum = (String) object.toString(); + String tsNum = object.toString(); timeseriesNumList = new ArrayList<>(); timeseriesNumList.add(tsNum); colCount = 1; @@ -87,7 +87,7 @@ public class IoTDBMetadataResultSet extends IoTDBQueryResultSet { columnItr = timeseriesNumList.iterator(); break; case COUNT_NODES: - String ndNum = (String) object.toString(); + String ndNum = object.toString(); nodesNumList = new ArrayList<>(); nodesNumList.add(ndNum); colCount = 1; diff --git a/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java b/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java index 3e36516..3bbc84f 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java @@ -37,6 +37,7 @@ import org.apache.iotdb.db.exception.ProcessorException; import org.apache.iotdb.db.exception.StorageEngineException; import org.apache.iotdb.db.exception.StorageEngineFailureException; import org.apache.iotdb.db.metadata.MManager; +import org.apache.iotdb.db.metadata.MNode; import org.apache.iotdb.db.qp.physical.crud.BatchInsertPlan; import org.apache.iotdb.db.qp.physical.crud.InsertPlan; import org.apache.iotdb.db.query.context.QueryContext; @@ -87,13 +88,14 @@ public class StorageEngine implements IService { * recover all storage group processors. */ try { - List<String> storageGroups = MManager.getInstance().getAllStorageGroupNames(); - for (String storageGroup : storageGroups) { - StorageGroupProcessor processor = new StorageGroupProcessor(systemDir, storageGroup); - logger.info("Storage Group Processor {} is recovered successfully", storageGroup); - processorMap.put(storageGroup, processor); + List<MNode> sgNodes = MManager.getInstance().getAllStorageGroups(); + for (MNode storageGroup : sgNodes) { + StorageGroupProcessor processor = new StorageGroupProcessor(systemDir, storageGroup.getFullPath()); + processor.setDataTTL(storageGroup.getDataTTL()); + logger.info("Storage Group Processor {} is recovered successfully", storageGroup.getFullPath()); + processorMap.put(storageGroup.getFullPath(), processor); } - } catch (ProcessorException | MetadataErrorException e) { + } catch (ProcessorException e) { logger.error("init a storage group processor failed. ", e); throw new StorageEngineFailureException(e); } @@ -326,13 +328,8 @@ public class StorageEngine implements IService { */ public synchronized boolean deleteAll() { logger.info("Start deleting all storage groups' timeseries"); - try { - for (String storageGroup : MManager.getInstance().getAllStorageGroupNames()) { - this.deleteAllDataFilesInOneStorageGroup(storageGroup); - } - } catch (MetadataErrorException e) { - logger.error("delete storage groups failed.", e); - return false; + for (String storageGroup : MManager.getInstance().getAllStorageGroupNames()) { + this.deleteAllDataFilesInOneStorageGroup(storageGroup); } return true; } diff --git a/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java b/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java index 2aa882e..f4690fd 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java @@ -144,12 +144,12 @@ public abstract class AbstractMemTable implements IMemTable { @Override public ReadOnlyMemChunk query(String deviceId, String measurement, TSDataType dataType, - Map<String, String> props) { + Map<String, String> props, long timeLowerBound) { TimeValuePairSorter sorter; if (!checkPath(deviceId, measurement)) { return null; } else { - long undeletedTime = findUndeletedTime(deviceId, measurement); + long undeletedTime = findUndeletedTime(deviceId, measurement, timeLowerBound); IWritableMemChunk memChunk = memTableMap.get(deviceId).get(measurement); IWritableMemChunk chunkCopy = new WritableMemChunk(dataType, memChunk.getTVList().clone()); chunkCopy.setTimeOffset(undeletedTime); @@ -159,7 +159,7 @@ public abstract class AbstractMemTable implements IMemTable { } - private long findUndeletedTime(String deviceId, String measurement) { + private long findUndeletedTime(String deviceId, String measurement, long timeLowerBound) { long undeletedTime = Long.MIN_VALUE; for (Modification modification : modifications) { if (modification instanceof Deletion) { @@ -170,7 +170,7 @@ public abstract class AbstractMemTable implements IMemTable { } } } - return undeletedTime + 1; + return undeletedTime + 1 < timeLowerBound ? timeLowerBound : undeletedTime + 1; } @Override diff --git a/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java b/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java index 5501e33..2e43380 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java @@ -57,7 +57,7 @@ public interface IMemTable { void insertBatch(BatchInsertPlan batchInsertPlan, List<Integer> indexes); ReadOnlyMemChunk query(String deviceId, String measurement, TSDataType dataType, - Map<String, String> props); + Map<String, String> props, long timeLowerBound); /** * putBack all the memory resources. diff --git a/server/src/main/java/org/apache/iotdb/db/engine/querycontext/QueryDataSource.java b/server/src/main/java/org/apache/iotdb/db/engine/querycontext/QueryDataSource.java index e5c87dd..c0e49f9 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/querycontext/QueryDataSource.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/querycontext/QueryDataSource.java @@ -23,12 +23,20 @@ import org.apache.iotdb.db.engine.storagegroup.TsFileResource; import org.apache.iotdb.tsfile.read.common.Path; import java.util.List; +import org.apache.iotdb.tsfile.read.filter.TimeFilter; +import org.apache.iotdb.tsfile.read.filter.basic.Filter; +import org.apache.iotdb.tsfile.read.filter.operator.AndFilter; public class QueryDataSource { private Path seriesPath; private List<TsFileResource> seqResources; private List<TsFileResource> unseqResources; + /** + * data older than currentTime - dataTTL should be ignored. + */ + private long dataTTL = Long.MAX_VALUE; + public QueryDataSource(Path seriesPath, List<TsFileResource> seqResources, List<TsFileResource> unseqResources) { this.seriesPath = seriesPath; this.seqResources = seqResources; @@ -46,4 +54,28 @@ public class QueryDataSource { public List<TsFileResource> getUnseqResources() { return unseqResources; } + + public long getDataTTL() { + return dataTTL; + } + + public void setDataTTL(long dataTTL) { + this.dataTTL = dataTTL; + } + + /** + * + * @return an updated time filter considering TTL + */ + public Filter updateTimeFilter(Filter timeFilter) { + if (dataTTL != Long.MAX_VALUE) { + if (timeFilter != null) { + timeFilter = new AndFilter(timeFilter, TimeFilter.gtEq(System.currentTimeMillis() - + dataTTL)); + } else { + timeFilter = TimeFilter.gtEq(System.currentTimeMillis() - dataTTL); + } + } + return timeFilter; + } } diff --git a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java index 389de7e..2ce3d42 100755 --- a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java @@ -114,7 +114,7 @@ public class StorageGroupProcessor { */ private final ReadWriteLock insertLock = new ReentrantReadWriteLock(); /** - * + * closeStorageGroupCondition is used to wait for all currently closing TsFiles to be done. */ private final Object closeStorageGroupCondition = new Object(); /** @@ -176,6 +176,11 @@ public class StorageGroupProcessor { private LinkedList<String> lruForSensorUsedInQuery = new LinkedList<>(); private static final int MAX_CACHE_SENSORS = 5000; + /** + * when the data in a storage group are older than dataTTL, it is considered invalid and will + * be eventually removed. + */ + private long dataTTL = Long.MAX_VALUE; public StorageGroupProcessor(String systemInfoDir, String storageGroupName) throws ProcessorException { @@ -299,7 +304,7 @@ public class StorageGroupProcessor { // TsFileNameComparator compares TsFiles by the version number in its name // ({systemTime}-{versionNum}.tsfile) - public int compareFileName(File o1, File o2) { + private int compareFileName(File o1, File o2) { String[] items1 = o1.getName().replace(TSFILE_SUFFIX, "").split("-"); String[] items2 = o2.getName().replace(TSFILE_SUFFIX, "").split("-"); if (Long.valueOf(items1[0]) - Long.valueOf(items2[0]) == 0) { @@ -336,6 +341,10 @@ public class StorageGroupProcessor { } public boolean insert(InsertPlan insertPlan) { + // reject insertions that are out of ttl + if (!checkTTL(insertPlan.getTime())) { + return false; + } writeLock(); try { // init map @@ -363,8 +372,13 @@ public class StorageGroupProcessor { for (int i = 0; i < batchInsertPlan.getRowCount(); i++) { results[i] = TSStatusType.SUCCESS_STATUS.getStatusCode(); - if (batchInsertPlan.getTimes()[i] > latestFlushedTimeForEachDevice - .get(batchInsertPlan.getDeviceId())) { + long currTime = batchInsertPlan.getTimes()[i]; + // skip points that do not satisfy TTL + if (!checkTTL(currTime)) { + results[i] = TSStatusType.OUT_OF_TTL_ERROR.getStatusCode(); + continue; + } + if (currTime > latestFlushedTimeForEachDevice.get(batchInsertPlan.getDeviceId())) { sequenceIndexes.add(i); } else { unsequenceIndexes.add(i); @@ -384,6 +398,15 @@ public class StorageGroupProcessor { } } + /** + * + * @param time + * @return whether the given time falls in ttl + */ + private boolean checkTTL(long time) { + return dataTTL == Long.MAX_VALUE || (System.currentTimeMillis() - time) <= dataTTL; + } + private void insertBatchToTsFileProcessor(BatchInsertPlan batchInsertPlan, List<Integer> indexes, boolean sequence, Integer[] results) { @@ -616,6 +639,7 @@ public class StorageGroupProcessor { if (filePathsManager != null) { filePathsManager.addUsedFilesForGivenJob(context.getJobId(), dataSource); } + dataSource.setDataTTL(dataTTL); return dataSource; } finally { insertLock.readLock().unlock(); @@ -664,31 +688,39 @@ public class StorageGroupProcessor { TSDataType dataType = mSchema.getType(); List<TsFileResource> tsfileResourcesForQuery = new ArrayList<>(); + long timeLowerBound = dataTTL != Long.MAX_VALUE ? System.currentTimeMillis() - dataTTL : Long + .MIN_VALUE; for (TsFileResource tsFileResource : tsFileResources) { // TODO: try filtering files if the query contains time filter if (!tsFileResource.containsDevice(deviceId)) { continue; } - if (!tsFileResource.getStartTimeMap().isEmpty()) { - closeQueryLock.readLock().lock(); - try { - if (tsFileResource.isClosed()) { - tsfileResourcesForQuery.add(tsFileResource); - } else { - // left: in-memory data, right: meta of disk data - Pair<ReadOnlyMemChunk, List<ChunkMetaData>> pair; - pair = tsFileResource - .getUnsealedFileProcessor() - .query(deviceId, measurementId, dataType, mSchema.getProps(), context); - tsfileResourcesForQuery - .add(new TsFileResource(tsFileResource.getFile(), - tsFileResource.getStartTimeMap(), - tsFileResource.getEndTimeMap(), pair.left, pair.right)); - } - } finally { - closeQueryLock.readLock().unlock(); + closeQueryLock.readLock().lock(); + + if (dataTTL != Long.MAX_VALUE) { + Long deviceEndTime = tsFileResource.getEndTimeMap().get(deviceId); + if (deviceEndTime != null && !checkTTL(deviceEndTime)) { + continue; } } + + try { + if (tsFileResource.isClosed()) { + tsfileResourcesForQuery.add(tsFileResource); + } else { + // left: in-memory data, right: meta of disk data + Pair<ReadOnlyMemChunk, List<ChunkMetaData>> pair; + pair = tsFileResource + .getUnsealedFileProcessor() + .query(deviceId, measurementId, dataType, mSchema.getProps(), context, timeLowerBound); + tsfileResourcesForQuery + .add(new TsFileResource(tsFileResource.getFile(), + tsFileResource.getStartTimeMap(), + tsFileResource.getEndTimeMap(), pair.left, pair.right)); + } + } finally { + closeQueryLock.readLock().unlock(); + } } return tsfileResourcesForQuery; } @@ -962,4 +994,11 @@ public class StorageGroupProcessor { void call(TsFileProcessor caller) throws TsFileProcessorException, IOException; } + public long getDataTTL() { + return dataTTL; + } + + public void setDataTTL(long dataTTL) { + this.dataTTL = dataTTL; + } } diff --git a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java index 5d237d1..7bc52a8 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java @@ -536,10 +536,12 @@ public class TsFileProcessor { * @param deviceId device id * @param measurementId sensor id * @param dataType data type + * @param timeLowerBound data older than this is ignored * @return left: the chunk data in memory; right: the chunkMetadatas of data on disk */ public Pair<ReadOnlyMemChunk, List<ChunkMetaData>> query(String deviceId, - String measurementId, TSDataType dataType, Map<String, String> props, QueryContext context) { + String measurementId, TSDataType dataType, Map<String, String> props, QueryContext context, + long timeLowerBound) { flushQueryLock.readLock().lock(); try { MemSeriesLazyMerger memSeriesLazyMerger = new MemSeriesLazyMerger(); @@ -548,13 +550,14 @@ public class TsFileProcessor { continue; } ReadOnlyMemChunk memChunk = flushingMemTable - .query(deviceId, measurementId, dataType, props); + .query(deviceId, measurementId, dataType, props, timeLowerBound); if (memChunk != null) { memSeriesLazyMerger.addMemSeries(memChunk); } } if (workMemTable != null) { - ReadOnlyMemChunk memChunk = workMemTable.query(deviceId, measurementId, dataType, props); + ReadOnlyMemChunk memChunk = workMemTable.query(deviceId, measurementId, dataType, props, + timeLowerBound); if (memChunk != null) { memSeriesLazyMerger.addMemSeries(memChunk); } @@ -573,6 +576,8 @@ public class TsFileProcessor { QueryUtils.modifyChunkMetaData(chunkMetaDataList, modifications); + chunkMetaDataList.removeIf(chunkMetaData -> chunkMetaData.getEndTime() < timeLowerBound); + return new Pair<>(timeValuePairSorter, chunkMetaDataList); } finally { flushQueryLock.readLock().unlock(); diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/MGraph.java b/server/src/main/java/org/apache/iotdb/db/metadata/MGraph.java index c35e710..6310628 100644 --- a/server/src/main/java/org/apache/iotdb/db/metadata/MGraph.java +++ b/server/src/main/java/org/apache/iotdb/db/metadata/MGraph.java @@ -177,7 +177,7 @@ public class MGraph implements Serializable { * * @return A HashMap whose Keys are separated by the storage file name. */ - HashMap<String, ArrayList<String>> getAllPathGroupByFilename(String path) + HashMap<String, ArrayList<String>> getAllPathGroupByStorageGroup(String path) throws PathErrorException { String rootName = path.trim().split(DOUB_SEPARATOR)[0]; if (mtree.getRoot().getName().equals(rootName)) { @@ -189,6 +189,10 @@ public class MGraph implements Serializable { throw new PathErrorException(TIME_SERIES_INCORRECT + rootName); } + List<MNode> getAllStorageGroups() { + return mtree.getAllStorageGroups(); + } + /** * function for getting all timeseries paths under the given seriesPath. */ @@ -243,7 +247,7 @@ public class MGraph implements Serializable { return new Metadata(deviceIdMap); } - HashSet<String> getAllStorageGroup() { + List<String> getAllStorageGroup() { return mtree.getAllStorageGroup(); } @@ -304,14 +308,14 @@ public class MGraph implements Serializable { return mtree.getStorageGroupNameByPath(node, path); } - boolean checkFileNameByPath(String path) { + boolean checkStorageGroupByPath(String path) { return mtree.checkFileNameByPath(path); } /** * Get all file names for given seriesPath */ - List<String> getAllFileNamesByPath(String path) throws PathErrorException { + List<String> getAllStorageGroupNamesByPath(String path) throws PathErrorException { return mtree.getAllFileNamesByPath(path); } @@ -384,7 +388,7 @@ public class MGraph implements Serializable { */ Map<String, Integer> countSeriesNumberInEachStorageGroup() throws PathErrorException { Map<String, Integer> res = new HashMap<>(); - Set<String> storageGroups = this.getAllStorageGroup(); + List<String> storageGroups = this.getAllStorageGroup(); for (String sg : storageGroups) { MNode node = mtree.getNodeByPath(sg); res.put(sg, node.getLeafCount()); diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java b/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java index 2da2f6c..e0b4d47 100644 --- a/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java +++ b/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java @@ -272,7 +272,7 @@ public class MManager { throw new MetadataErrorException( String.format("Timeseries %s already exist", path.getFullPath())); } - if (!checkFileNameByPath(path.getFullPath())) { + if (!checkStorageGroupByPath(path.getFullPath())) { throw new MetadataErrorException("Storage group should be created first"); } // optimize the speed of adding timeseries @@ -489,7 +489,7 @@ public class MManager { try { checkAndGetDataTypeCache.clear(); mNodeCache.clear(); - String dataFileName = mgraph.deletePath(path); + String storageGroupName = mgraph.deletePath(path); if (writeToLog) { BufferedWriter writer = getLogWriter(); writer.write(MetadataOperationType.DELETE_PATH_FROM_MTREE + "," + path); @@ -510,7 +510,7 @@ public class MManager { } else { maxSeriesNumberAmongStorageGroup--; } - return dataFileName; + return storageGroupName; } finally { lock.writeLock().unlock(); } @@ -745,21 +745,6 @@ public class MManager { } /** - * Get the full storage group info. - * - * @return A HashSet instance which stores all storage group info - */ - public Set<String> getAllStorageGroup() throws PathErrorException { - - lock.readLock().lock(); - try { - return mgraph.getAllStorageGroup(); - } finally { - lock.readLock().unlock(); - } - } - - /** * Get all nodes from the given level * * @return A List instance which stores all node at given level @@ -877,42 +862,55 @@ public class MManager { } /** - * function for checking file name by path. + * function for checking storage group name by path. */ - boolean checkFileNameByPath(String path) { + boolean checkStorageGroupByPath(String path) { lock.readLock().lock(); try { - return mgraph.checkFileNameByPath(path); + return mgraph.checkStorageGroupByPath(path); } finally { lock.readLock().unlock(); } } /** - * function for getting all file names. + * Get the full storage group info. + * + * @return A list which stores all storage group info */ - public List<String> getAllStorageGroupNames() throws MetadataErrorException { + public List<String> getAllStorageGroupNames() { + + lock.readLock().lock(); + try { + return mgraph.getAllStorageGroup(); + } finally { + lock.readLock().unlock(); + } + } + /** + * function for getting all storage groups' MNodes + */ + public List<MNode> getAllStorageGroups() { lock.readLock().lock(); try { - Map<String, ArrayList<String>> res = getAllPathGroupByFileName(ROOT_NAME); - return new ArrayList<>(res.keySet()); + return mgraph.getAllStorageGroups(); } finally { lock.readLock().unlock(); } } /** - * Get all file names for given seriesPath + * Get all storage group names for given seriesPath * - * @return List of String represented all file names + * @return List of String represented all storage group names */ - List<String> getAllFileNamesByPath(String path) throws MetadataErrorException { + List<String> getAllStorageGroupNamesByPath(String path) throws MetadataErrorException { lock.readLock().lock(); try { - return mgraph.getAllFileNamesByPath(path); + return mgraph.getAllStorageGroupNamesByPath(path); } catch (PathErrorException e) { throw new MetadataErrorException(e); } finally { @@ -921,13 +919,13 @@ public class MManager { } /** - * return a HashMap contains all the paths separated by File Name. + * return a HashMap contains all the paths separated by storage group name. */ - Map<String, ArrayList<String>> getAllPathGroupByFileName(String path) + Map<String, ArrayList<String>> getAllPathGroupByStorageGroup(String path) throws MetadataErrorException { lock.readLock().lock(); try { - return mgraph.getAllPathGroupByFilename(path); + return mgraph.getAllPathGroupByStorageGroup(path); } catch (PathErrorException e) { throw new MetadataErrorException(e); } finally { @@ -944,8 +942,8 @@ public class MManager { lock.readLock().lock(); try { ArrayList<String> res = new ArrayList<>(); - Map<String, ArrayList<String>> pathsGroupByFilename = getAllPathGroupByFileName(path); - for (ArrayList<String> ps : pathsGroupByFilename.values()) { + Map<String, ArrayList<String>> pathsGroupBySG = getAllPathGroupByStorageGroup(path); + for (ArrayList<String> ps : pathsGroupBySG.values()) { res.addAll(ps); } return res; diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/MNode.java b/server/src/main/java/org/apache/iotdb/db/metadata/MNode.java index 433b7d5..88f8662 100644 --- a/server/src/main/java/org/apache/iotdb/db/metadata/MNode.java +++ b/server/src/main/java/org/apache/iotdb/db/metadata/MNode.java @@ -22,6 +22,7 @@ import java.io.Serializable; import java.util.HashMap; import java.util.LinkedHashMap; import java.util.Map; +import org.apache.iotdb.db.conf.IoTDBConstant; import org.apache.iotdb.tsfile.file.metadata.enums.CompressionType; import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType; import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding; @@ -52,6 +53,15 @@ public class MNode implements Serializable { private MNode parent; private Map<String, MNode> children; + private String fullPath; + + /** + * when the data in a storage group are older than dataTTL, it is considered invalid and will + * be eventually removed. + * only at storage group will this be set. + */ + private long dataTTL = Long.MAX_VALUE; + /** * Constructor of MNode. */ @@ -208,4 +218,24 @@ public class MNode implements Serializable { this.name = name; } + public long getDataTTL() { + return dataTTL; + } + + public void setDataTTL(long dataTTL) { + this.dataTTL = dataTTL; + } + + public String getFullPath() { + if (fullPath != null) { + return fullPath; + } + StringBuilder builder = new StringBuilder(name); + MNode curr = this; + while (curr.parent != null) { + curr = curr.parent; + builder.insert(0, IoTDBConstant.PATH_SEPARATOR).insert(0, curr.name); + } + return fullPath = builder.toString(); + } } diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/MTree.java b/server/src/main/java/org/apache/iotdb/db/metadata/MTree.java index df0032b..69acdc6 100644 --- a/server/src/main/java/org/apache/iotdb/db/metadata/MTree.java +++ b/server/src/main/java/org/apache/iotdb/db/metadata/MTree.java @@ -652,6 +652,26 @@ public class MTree implements Serializable { } /** + * + * @return all storage groups' MNodes + * @throws PathErrorException + */ + List<MNode> getAllStorageGroups() { + List<MNode> ret = new ArrayList<>(); + Stack<MNode> nodeStack = new Stack<>(); + nodeStack.add(getRoot()); + while (!nodeStack.isEmpty()) { + MNode current = nodeStack.pop(); + if (current.isStorageLevel()) { + ret.add(current); + } else if (current.hasChildren()){ + nodeStack.addAll(current.getChildren().values()); + } + } + return ret; + } + + /** * function for getting all timeseries paths under the given seriesPath. */ List<List<String>> getShowTimeseriesPath(String pathReg) throws PathErrorException { @@ -740,8 +760,8 @@ public class MTree implements Serializable { * * @return a list contains all distinct storage groups */ - HashSet<String> getAllStorageGroup() { - HashSet<String> res = new HashSet<>(); + List<String> getAllStorageGroup() { + List<String> res = new ArrayList<>(); MNode rootNode; if ((rootNode = getRoot()) != null) { findStorageGroup(rootNode, "root", res); @@ -749,7 +769,7 @@ public class MTree implements Serializable { return res; } - private void findStorageGroup(MNode node, String path, HashSet<String> res) { + private void findStorageGroup(MNode node, String path, List<String> res) { if (node.isStorageLevel()) { res.add(path); return; diff --git a/server/src/main/java/org/apache/iotdb/db/query/dataset/groupby/GroupByWithoutValueFilterDataSet.java b/server/src/main/java/org/apache/iotdb/db/query/dataset/groupby/GroupByWithoutValueFilterDataSet.java index 3043d09..9c3b873 100644 --- a/server/src/main/java/org/apache/iotdb/db/query/dataset/groupby/GroupByWithoutValueFilterDataSet.java +++ b/server/src/main/java/org/apache/iotdb/db/query/dataset/groupby/GroupByWithoutValueFilterDataSet.java @@ -82,6 +82,7 @@ public class GroupByWithoutValueFilterDataSet extends GroupByEngineDataSet { for (Path path : selectedSeries) { QueryDataSource queryDataSource = QueryResourceManager.getInstance() .getQueryDataSource(path, context); + timeFilter = queryDataSource.updateTimeFilter(timeFilter); // sequence reader for sealed tsfile, unsealed tsfile, memory IAggregateReader seqResourceIterateReader = new SeqResourceIterateReader( diff --git a/server/src/main/java/org/apache/iotdb/db/query/executor/AggregateEngineExecutor.java b/server/src/main/java/org/apache/iotdb/db/query/executor/AggregateEngineExecutor.java index 994ffa6..a69d204 100644 --- a/server/src/main/java/org/apache/iotdb/db/query/executor/AggregateEngineExecutor.java +++ b/server/src/main/java/org/apache/iotdb/db/query/executor/AggregateEngineExecutor.java @@ -50,7 +50,9 @@ import org.apache.iotdb.tsfile.read.common.BatchData; import org.apache.iotdb.tsfile.read.common.Path; import org.apache.iotdb.tsfile.read.expression.IExpression; import org.apache.iotdb.tsfile.read.expression.impl.GlobalTimeExpression; +import org.apache.iotdb.tsfile.read.filter.TimeFilter; import org.apache.iotdb.tsfile.read.filter.basic.Filter; +import org.apache.iotdb.tsfile.read.filter.operator.AndFilter; import org.apache.iotdb.tsfile.read.query.dataset.QueryDataSet; public class AggregateEngineExecutor { @@ -100,6 +102,8 @@ public class AggregateEngineExecutor { QueryDataSource queryDataSource = QueryResourceManager.getInstance() .getQueryDataSource(selectedSeries.get(i), context); + // add additional time filter if TTL is set + timeFilter = queryDataSource.updateTimeFilter(timeFilter); // sequence reader for sealed tsfile, unsealed tsfile, memory IAggregateReader seqResourceIterateReader; diff --git a/server/src/main/java/org/apache/iotdb/db/query/reader/seriesRelated/SeriesReaderWithoutValueFilter.java b/server/src/main/java/org/apache/iotdb/db/query/reader/seriesRelated/SeriesReaderWithoutValueFilter.java index 901352f..1e0c207 100644 --- a/server/src/main/java/org/apache/iotdb/db/query/reader/seriesRelated/SeriesReaderWithoutValueFilter.java +++ b/server/src/main/java/org/apache/iotdb/db/query/reader/seriesRelated/SeriesReaderWithoutValueFilter.java @@ -64,25 +64,26 @@ public class SeriesReaderWithoutValueFilter implements IPointReader { * Constructor function. * * @param seriesPath the path of the series data - * @param filter filter condition + * @param timeFilter time filter condition * @param context query context * @param pushdownUnseq True to push down the filter on the unsequence TsFile resource; False not * to. */ - protected SeriesReaderWithoutValueFilter(Path seriesPath, Filter filter, QueryContext context, + protected SeriesReaderWithoutValueFilter(Path seriesPath, Filter timeFilter, QueryContext context, boolean pushdownUnseq) throws StorageEngineException, IOException { QueryDataSource queryDataSource = QueryResourceManager.getInstance() .getQueryDataSource(seriesPath, context); + timeFilter = queryDataSource.updateTimeFilter(timeFilter); // reader for sequence resources IBatchReader seqResourceIterateReader = new SeqResourceIterateReader( - queryDataSource.getSeriesPath(), queryDataSource.getSeqResources(), filter, context); + queryDataSource.getSeriesPath(), queryDataSource.getSeqResources(), timeFilter, context); // reader for unsequence resources IPointReader unseqResourceMergeReader; if (pushdownUnseq) { unseqResourceMergeReader = new UnseqResourceMergeReader(seriesPath, - queryDataSource.getUnseqResources(), context, filter); + queryDataSource.getUnseqResources(), context, timeFilter); } else { unseqResourceMergeReader = new UnseqResourceMergeReader(seriesPath, queryDataSource.getUnseqResources(), context, null); diff --git a/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java b/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java index ae6496b..4368cd2 100644 --- a/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java +++ b/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java @@ -75,18 +75,6 @@ import org.apache.thrift.server.ServerContext; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.io.IOException; -import java.nio.ByteBuffer; -import java.sql.SQLException; -import java.sql.Statement; -import java.time.ZoneId; -import java.util.*; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.atomic.AtomicLong; -import java.util.regex.Pattern; - -import static org.apache.iotdb.db.conf.IoTDBConstant.*; - /** * Thrift RPC implementation at server side. */ @@ -271,7 +259,7 @@ public class TSServiceImpl implements TSIService.Iface, ServerContext { status = new TS_Status(getStatus(TSStatusType.SUCCESS_STATUS)); break; case "SHOW_STORAGE_GROUP": - Set<String> storageGroups = getAllStorageGroups(); + List<String> storageGroups = getAllStorageGroups(); resp.setShowStorageGroups(storageGroups); status = new TS_Status(getStatus(TSStatusType.SUCCESS_STATUS)); break; @@ -337,8 +325,8 @@ public class TSServiceImpl implements TSIService.Iface, ServerContext { return MManager.getInstance().getNodesList(level); } - private Set<String> getAllStorageGroups() throws PathErrorException { - return MManager.getInstance().getAllStorageGroup(); + private List<String> getAllStorageGroups() throws PathErrorException { + return MManager.getInstance().getAllStorageGroupNames(); } private List<List<String>> getTimeSeriesForPath(String path) @@ -526,7 +514,7 @@ public class TSServiceImpl implements TSIService.Iface, ServerContext { IoTDBDescriptor.getInstance().getConfig().getMaxMemtableNumber(), IoTDBDescriptor.getInstance().getConfig().getTsFileSizeThreshold(), CompressionRatio.getInstance().getRatio(), - MManager.getInstance().getAllStorageGroup().size(), + MManager.getInstance().getAllStorageGroupNames().size(), IoTDBConfigDynamicAdapter.getInstance().getTotalTimeseries(), MManager.getInstance().getMaximalSeriesNumberAmongStorageGroups()); return getTSExecuteStatementResp(getStatus(TSStatusType.SUCCESS_STATUS, msg)); diff --git a/server/src/test/java/org/apache/iotdb/db/metadata/MManagerAdvancedTest.java b/server/src/test/java/org/apache/iotdb/db/metadata/MManagerAdvancedTest.java index 94bfd40..de4e9ec 100644 --- a/server/src/test/java/org/apache/iotdb/db/metadata/MManagerAdvancedTest.java +++ b/server/src/test/java/org/apache/iotdb/db/metadata/MManagerAdvancedTest.java @@ -84,7 +84,7 @@ public class MManagerAdvancedTest { // test filename by seriesPath assertEquals("root.vehicle.d0", mmanager.getStorageGroupNameByPath("root.vehicle.d0.s1")); Map<String, ArrayList<String>> map = mmanager - .getAllPathGroupByFileName("root.vehicle.d1.*"); + .getAllPathGroupByStorageGroup("root.vehicle.d1.*"); assertEquals(1, map.keySet().size()); assertEquals(6, map.get("root.vehicle.d1").size()); List<String> paths = mmanager.getPaths("root.vehicle.d0"); diff --git a/server/src/test/java/org/apache/iotdb/db/metadata/MManagerBasicTest.java b/server/src/test/java/org/apache/iotdb/db/metadata/MManagerBasicTest.java index 864381c..866a5fb 100644 --- a/server/src/test/java/org/apache/iotdb/db/metadata/MManagerBasicTest.java +++ b/server/src/test/java/org/apache/iotdb/db/metadata/MManagerBasicTest.java @@ -135,7 +135,7 @@ public class MManagerBasicTest { } assertFalse(manager.pathExist("root.laptop.d2")); - assertFalse(manager.checkFileNameByPath("root.laptop.d2")); + assertFalse(manager.checkStorageGroupByPath("root.laptop.d2")); try { manager.deletePaths(Collections.singletonList(new Path("root.laptop.d1.s0"))); @@ -299,12 +299,12 @@ public class MManagerBasicTest { List<String> list = new ArrayList<>(); list.add("root.laptop.d1"); - assertEquals(list, manager.getAllFileNamesByPath("root.laptop.d1.s1")); - assertEquals(list, manager.getAllFileNamesByPath("root.laptop.d1")); + assertEquals(list, manager.getAllStorageGroupNamesByPath("root.laptop.d1.s1")); + assertEquals(list, manager.getAllStorageGroupNamesByPath("root.laptop.d1")); list.add("root.laptop.d2"); - assertEquals(list, manager.getAllFileNamesByPath("root.laptop")); - assertEquals(list, manager.getAllFileNamesByPath("root")); + assertEquals(list, manager.getAllStorageGroupNamesByPath("root.laptop")); + assertEquals(list, manager.getAllStorageGroupNamesByPath("root")); } catch (MetadataErrorException e) { e.printStackTrace(); fail(e.getMessage()); @@ -316,23 +316,23 @@ public class MManagerBasicTest { MManager manager = MManager.getInstance(); try { - assertTrue(manager.getAllPathGroupByFileName("root").keySet().isEmpty()); - assertTrue(manager.getAllFileNamesByPath("root.vehicle").isEmpty()); - assertTrue(manager.getAllFileNamesByPath("root.vehicle.device").isEmpty()); - assertTrue(manager.getAllFileNamesByPath("root.vehicle.device.sensor").isEmpty()); + assertTrue(manager.getAllPathGroupByStorageGroup("root").keySet().isEmpty()); + assertTrue(manager.getAllStorageGroupNamesByPath("root.vehicle").isEmpty()); + assertTrue(manager.getAllStorageGroupNamesByPath("root.vehicle.device").isEmpty()); + assertTrue(manager.getAllStorageGroupNamesByPath("root.vehicle.device.sensor").isEmpty()); manager.setStorageLevelToMTree("root.vehicle"); - assertFalse(manager.getAllFileNamesByPath("root.vehicle").isEmpty()); - assertFalse(manager.getAllFileNamesByPath("root.vehicle.device").isEmpty()); - assertFalse(manager.getAllFileNamesByPath("root.vehicle.device.sensor").isEmpty()); - assertTrue(manager.getAllFileNamesByPath("root.vehicle1").isEmpty()); - assertTrue(manager.getAllFileNamesByPath("root.vehicle1.device").isEmpty()); + assertFalse(manager.getAllStorageGroupNamesByPath("root.vehicle").isEmpty()); + assertFalse(manager.getAllStorageGroupNamesByPath("root.vehicle.device").isEmpty()); + assertFalse(manager.getAllStorageGroupNamesByPath("root.vehicle.device.sensor").isEmpty()); + assertTrue(manager.getAllStorageGroupNamesByPath("root.vehicle1").isEmpty()); + assertTrue(manager.getAllStorageGroupNamesByPath("root.vehicle1.device").isEmpty()); manager.setStorageLevelToMTree("root.vehicle1.device"); - assertTrue(manager.getAllFileNamesByPath("root.vehicle1.device1").isEmpty()); - assertTrue(manager.getAllFileNamesByPath("root.vehicle1.device2").isEmpty()); - assertTrue(manager.getAllFileNamesByPath("root.vehicle1.device3").isEmpty()); - assertFalse(manager.getAllFileNamesByPath("root.vehicle1.device").isEmpty()); + assertTrue(manager.getAllStorageGroupNamesByPath("root.vehicle1.device1").isEmpty()); + assertTrue(manager.getAllStorageGroupNamesByPath("root.vehicle1.device2").isEmpty()); + assertTrue(manager.getAllStorageGroupNamesByPath("root.vehicle1.device3").isEmpty()); + assertFalse(manager.getAllStorageGroupNamesByPath("root.vehicle1.device").isEmpty()); } catch (MetadataErrorException e) { e.printStackTrace(); fail(e.getMessage()); diff --git a/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusType.java b/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusType.java index 5a4acf6..58c8b4e 100644 --- a/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusType.java +++ b/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusType.java @@ -27,6 +27,7 @@ public enum TSStatusType { UNSUPPORTED_FETCH_METADATA_OPERATION_ERROR(302, "Unsupported fetch metadata operation"), FETCH_METADATA_ERROR(303, "Failed to fetch metadata"), CHECK_FILE_LEVEL_ERROR(304, "Meet error while checking file level"), + OUT_OF_TTL_ERROR(305, "timestamp falls out of TTL"), EXECUTE_STATEMENT_ERROR(400, "Execute statement error"), SQL_PARSE_ERROR(401, "Meet error while parsing SQL"), GENERATE_TIME_ZONE_ERROR(402, "Meet error while generating time zone"), diff --git a/service-rpc/src/main/thrift/rpc.thrift b/service-rpc/src/main/thrift/rpc.thrift index 5cedb0c..ff2922b 100644 --- a/service-rpc/src/main/thrift/rpc.thrift +++ b/service-rpc/src/main/thrift/rpc.thrift @@ -195,7 +195,7 @@ struct TSFetchMetadataResp{ 3: optional list<string> ColumnsList 4: optional string dataType 5: optional list<list<string>> showTimeseriesList - 7: optional set<string> showStorageGroups + 7: optional list<string> showStorageGroups 8: optional list<string> nodesList 9: optional map<string, string> nodeTimeseriesNum }
