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;
+  }
 }

Reply via email to