This is an automated email from the ASF dual-hosted git repository. haonan pushed a commit to branch jira-1355 in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit a304b014d55bd37086ac0ff21851247e4b46e0a8 Author: HTHou <[email protected]> AuthorDate: Thu May 6 11:33:06 2021 +0800 [IOTDB-1355] Support updating aligned timeseries values when insert partially --- .../iotdb/db/engine/flush/MemTableFlushTask.java | 69 ++++++++++++++++++++++ .../iotdb/db/utils/datastructure/VectorTVList.java | 10 ++++ 2 files changed, 79 insertions(+) diff --git a/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java b/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java index 633dd6e..e75bbce 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java @@ -39,6 +39,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.IOException; +import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.concurrent.ExecutionException; @@ -161,12 +162,80 @@ public class MemTableFlushTask { new Runnable() { private void writeOneSeries( TVList tvPairs, IChunkWriter seriesWriterImpl, TSDataType dataType) { + List<Integer> duplicatedSortedRowIndexList = null; for (int sortedRowIndex = 0; sortedRowIndex < tvPairs.size(); sortedRowIndex++) { long time = tvPairs.getTime(sortedRowIndex); // skip duplicated data if ((sortedRowIndex + 1 < tvPairs.size() && (time == tvPairs.getTime(sortedRowIndex + 1)))) { + // record the time duplicated row index list for vector type + if (dataType == TSDataType.VECTOR) { + if (duplicatedSortedRowIndexList == null) { + duplicatedSortedRowIndexList = new ArrayList<>(); + } + duplicatedSortedRowIndexList.add(sortedRowIndex); + } + continue; + } + + // For vector type only, combine the time duplicated vector rows to one row + if (duplicatedSortedRowIndexList != null && !duplicatedSortedRowIndexList.isEmpty()) { + VectorTVList vectorTvPairs = (VectorTVList) tvPairs; + List<TSDataType> dataTypes = vectorTvPairs.getTsDataTypes(); + List<Integer> originRowIndexList = new ArrayList<>(); + duplicatedSortedRowIndexList.forEach(i -> { + originRowIndexList.add(vectorTvPairs.getValueIndex(i)); + }); + for (int columnIndex = 0; columnIndex < dataTypes.size(); columnIndex++) { + int validOriginRowIndex = vectorTvPairs.getValidRowIndex(originRowIndexList, columnIndex); + boolean isNull = vectorTvPairs.isValueMarked(validOriginRowIndex, columnIndex); + switch (dataTypes.get(columnIndex)) { + case BOOLEAN: + seriesWriterImpl.write( + time, + vectorTvPairs.getBooleanByValueIndex(validOriginRowIndex, columnIndex), + isNull); + break; + case INT32: + seriesWriterImpl.write( + time, + vectorTvPairs.getIntByValueIndex(validOriginRowIndex, columnIndex), + isNull); + break; + case INT64: + seriesWriterImpl.write( + time, + vectorTvPairs.getLongByValueIndex(validOriginRowIndex, columnIndex), + isNull); + break; + case FLOAT: + seriesWriterImpl.write( + time, + vectorTvPairs.getFloatByValueIndex(validOriginRowIndex, columnIndex), + isNull); + break; + case DOUBLE: + seriesWriterImpl.write( + time, + vectorTvPairs.getDoubleByValueIndex(validOriginRowIndex, columnIndex), + isNull); + break; + case TEXT: + seriesWriterImpl.write( + time, + vectorTvPairs.getBinaryByValueIndex(validOriginRowIndex, columnIndex), + isNull); + break; + default: + LOGGER.error( + "Storage group {} does not support data type: {}", + storageGroup, + dataType); + break; + } + } + seriesWriterImpl.write(time); continue; } diff --git a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/VectorTVList.java b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/VectorTVList.java index 0a3b405..441b9ef 100644 --- a/server/src/main/java/org/apache/iotdb/db/utils/datastructure/VectorTVList.java +++ b/server/src/main/java/org/apache/iotdb/db/utils/datastructure/VectorTVList.java @@ -650,4 +650,14 @@ public class VectorTVList extends TVList { public TSDataType getDataType() { return TSDataType.VECTOR; } + + public int getValidRowIndex(List<Integer> originRowIndexList, int columnIndex) { + int validRowIndex = originRowIndexList.get(0); + for (int originRowIndex : originRowIndexList) { + if (!isValueMarked(originRowIndex, columnIndex)) { + validRowIndex = originRowIndex; + } + } + return validRowIndex; + } }
