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,
