This is an automated email from the ASF dual-hosted git repository.

haonan pushed a commit to branch fix_insertRows_overlap
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit 8848346c03cb8122c991f9858be5cad3e3abb912
Author: HTHou <[email protected]>
AuthorDate: Wed Jan 10 12:11:28 2024 +0800

    Fix insertRecords may cause sequence data overlapped
---
 .../db/storageengine/dataregion/DataRegion.java    | 10 ++++--
 .../storageengine/dataregion/DataRegionTest.java   | 38 ++++++++++++++++++++++
 2 files changed, 45 insertions(+), 3 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
index 9e6b22bad00..dc587a56a87 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
@@ -1174,6 +1174,7 @@ public class DataRegion implements IDataRegionForQuery {
       InsertRowsNode insertRowsNode, boolean[] areSequence, long[] 
timePartitionIds) {
     List<InsertRowNode> executedInsertRowNodeList = new ArrayList<>();
     long[] costsForMetrics = new long[4];
+    Map<TsFileProcessor, Boolean> tsFileProcessorMap = new HashMap<>();
     for (int i = 0; i < areSequence.length; i++) {
       InsertRowNode insertRowNode = 
insertRowsNode.getInsertRowNodeList().get(i);
       if (insertRowNode.allMeasurementFailed()) {
@@ -1184,16 +1185,19 @@ public class DataRegion implements IDataRegionForQuery {
       if (tsFileProcessor == null) {
         continue;
       }
+      tsFileProcessorMap.put(tsFileProcessor, areSequence[i]);
       try {
         tsFileProcessor.insert(insertRowNode, costsForMetrics);
       } catch (WriteProcessException e) {
         insertRowsNode.getResults().put(i, 
RpcUtils.getStatus(e.getErrorCode(), e.getMessage()));
       }
       executedInsertRowNodeList.add(insertRowNode);
+    }
 
-      // check memtable size and may asyncTryToFlush the work memtable
-      if (tsFileProcessor.shouldFlush()) {
-        fileFlushPolicy.apply(this, tsFileProcessor, areSequence[i]);
+    // check memtable size and may asyncTryToFlush the work memtable
+    for (Map.Entry<TsFileProcessor, Boolean> entry : 
tsFileProcessorMap.entrySet()) {
+      if (entry.getKey().shouldFlush()) {
+        fileFlushPolicy.apply(this, entry.getKey(), entry.getValue());
       }
     }
 
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/DataRegionTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/DataRegionTest.java
index 18bfeba6838..a0429827360 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/DataRegionTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/DataRegionTest.java
@@ -31,10 +31,13 @@ import org.apache.iotdb.db.conf.IoTDBConfig;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.exception.DataRegionException;
 import org.apache.iotdb.db.exception.WriteProcessException;
+import org.apache.iotdb.db.exception.WriteProcessRejectException;
 import org.apache.iotdb.db.exception.query.QueryProcessException;
 import org.apache.iotdb.db.queryengine.common.QueryId;
 import org.apache.iotdb.db.queryengine.execution.fragment.QueryContext;
+import org.apache.iotdb.db.queryengine.plan.planner.plan.node.PlanNodeId;
 import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowNode;
+import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowsNode;
 import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertTabletNode;
 import org.apache.iotdb.db.storageengine.StorageEngine;
 import 
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.performer.ICompactionPerformer;
@@ -859,6 +862,41 @@ public class DataRegionTest {
     COMMON_CONFIG.setTimePartitionInterval(defaultTimePartition);
   }
 
+  @Test
+  public void testInsertUnSequenceRows()
+      throws IllegalPathException, WriteProcessRejectException, 
QueryProcessException,
+          DataRegionException {
+    int defaultAvgSeriesPointNumberThreshold = 
config.getAvgSeriesPointNumberThreshold();
+    config.setAvgSeriesPointNumberThreshold(2);
+    DataRegion dataRegion1 = new DummyDataRegion(systemDir, "root.Rows");
+    long[] time = new long[] {3, 4, 1, 2};
+    List<Integer> indexList = new ArrayList<>();
+    List<InsertRowNode> nodes = new ArrayList<>();
+    for (int i = 0; i < 4; i++) {
+      TSRecord record = new TSRecord(time[i], "root.Rows");
+      record.addTuple(DataPoint.getDataPoint(TSDataType.INT32, measurementId, 
String.valueOf(i)));
+      nodes.add(buildInsertRowNodeByTSRecord(record));
+      indexList.add(i);
+    }
+    InsertRowsNode insertRowsNode = new InsertRowsNode(new PlanNodeId(""), 
indexList, nodes);
+    dataRegion1.insert(insertRowsNode);
+    dataRegion1.syncCloseAllWorkingTsFileProcessors();
+    QueryDataSource queryDataSource =
+        dataRegion1.query(
+            Collections.singletonList(new PartialPath("root.Rows", 
measurementId)),
+            "root.Rows",
+            context,
+            null,
+            null);
+    Assert.assertEquals(1, queryDataSource.getSeqResources().size());
+    Assert.assertEquals(0, queryDataSource.getUnseqResources().size());
+    for (TsFileResource resource : queryDataSource.getSeqResources()) {
+      Assert.assertTrue(resource.isClosed());
+    }
+    dataRegion1.syncDeleteDataFiles();
+    
config.setAvgSeriesPointNumberThreshold(defaultAvgSeriesPointNumberThreshold);
+  }
+
   @Test
   public void testSmallReportProportionInsertRow()
       throws WriteProcessException, QueryProcessException, 
IllegalPathException, IOException,

Reply via email to