This is an automated email from the ASF dual-hosted git repository.
haonan pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new a7f6d755482 Optimize InsertRecords performance by reducing record
metrics (#11723)
a7f6d755482 is described below
commit a7f6d755482438c132bd263a07d1bf018d8e75e4
Author: Haonan <[email protected]>
AuthorDate: Thu Jan 4 20:26:26 2024 +0800
Optimize InsertRecords performance by reducing record metrics (#11723)
---
.../dataregion/DataExecutionVisitor.java | 3 +
.../analyze/cache/schema/DataNodeSchemaCache.java | 22 +++
.../db/storageengine/dataregion/DataRegion.java | 179 +++++++++++++++++----
.../dataregion/memtable/TsFileProcessor.java | 16 +-
.../dataregion/memtable/TsFileProcessorTest.java | 16 +-
5 files changed, 194 insertions(+), 42 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/statemachine/dataregion/DataExecutionVisitor.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/statemachine/dataregion/DataExecutionVisitor.java
index c685a949356..30c151de683 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/statemachine/dataregion/DataExecutionVisitor.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/statemachine/dataregion/DataExecutionVisitor.java
@@ -112,6 +112,9 @@ public class DataExecutionVisitor extends
PlanVisitor<TSStatus, DataRegion> {
try {
dataRegion.insert(node);
return StatusUtils.OK;
+ } catch (WriteProcessRejectException e) {
+ LOGGER.warn("Reject in executing plan node: {}, caused by {}", node,
e.getMessage());
+ return RpcUtils.getStatus(e.getErrorCode(), e.getMessage());
} catch (BatchProcessException e) {
LOGGER.warn("Batch failure in executing a InsertRowsNode.");
TSStatus firstStatus = null;
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/cache/schema/DataNodeSchemaCache.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/cache/schema/DataNodeSchemaCache.java
index f8198e4216b..5559f7b5483 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/cache/schema/DataNodeSchemaCache.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/cache/schema/DataNodeSchemaCache.java
@@ -257,6 +257,28 @@ public class DataNodeSchemaCache {
}
}
+ public void updateLastCacheWithoutLock(
+ String database,
+ PartialPath devicePath,
+ String[] measurements,
+ MeasurementSchema[] measurementSchemas,
+ boolean isAligned,
+ IntFunction<TimeValuePair> timeValuePairProvider,
+ IntPredicate shouldUpdateProvider,
+ boolean highPriorityUpdate,
+ Long latestFlushedTime) {
+ timeSeriesSchemaCache.updateLastCache(
+ database,
+ devicePath,
+ measurements,
+ measurementSchemas,
+ isAligned,
+ timeValuePairProvider,
+ shouldUpdateProvider,
+ highPriorityUpdate,
+ latestFlushedTime);
+ }
+
/**
* get or create SchemaCacheEntry and update last cache, only support
non-aligned sensor or
* aligned sensor without only one sub sensor
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 84b61aca712..59d2c59b0d5 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
@@ -1003,7 +1003,7 @@ public class DataRegion implements IDataRegionForQuery {
&& noFailure;
}
startTime = System.nanoTime();
- tryToUpdateBatchInsertLastCache(insertTabletNode);
+ tryToUpdateInsertTabletLastCache(insertTabletNode);
PERFORMANCE_OVERVIEW_METRICS.recordScheduleUpdateLastCacheCost(System.nanoTime()
- startTime);
if (!noFailure) {
@@ -1076,7 +1076,7 @@ public class DataRegion implements IDataRegionForQuery {
return true;
}
- private void tryToUpdateBatchInsertLastCache(InsertTabletNode node) {
+ private void tryToUpdateInsertTabletLastCache(InsertTabletNode node) {
if (!CommonDescriptor.getInstance().getConfig().isLastCacheEnable()
||
(config.getDataRegionConsensusProtocolClass().equals(ConsensusFactory.IOT_CONSENSUS)
&& node.isSyncFromLeaderWhenUsingIoTConsensus())) {
@@ -1119,25 +1119,31 @@ public class DataRegion implements IDataRegionForQuery {
if (tsFileProcessor == null) {
return;
}
-
- tsFileProcessor.insert(insertRowNode);
- long startTime = System.nanoTime();
- tryToUpdateInsertLastCache(insertRowNode);
-
PERFORMANCE_OVERVIEW_METRICS.recordScheduleUpdateLastCacheCost(System.nanoTime()
- startTime);
+ long[] costsForMetrics = new long[4];
+ tsFileProcessor.insert(insertRowNode, costsForMetrics);
+
PERFORMANCE_OVERVIEW_METRICS.recordCreateMemtableBlockCost(costsForMetrics[0]);
+
PERFORMANCE_OVERVIEW_METRICS.recordScheduleMemoryBlockCost(costsForMetrics[1]);
+ PERFORMANCE_OVERVIEW_METRICS.recordScheduleWalCost(costsForMetrics[2]);
+
PERFORMANCE_OVERVIEW_METRICS.recordScheduleMemTableCost(costsForMetrics[3]);
// check memtable size and may asyncTryToFlush the work memtable
if (tsFileProcessor.shouldFlush()) {
fileFlushPolicy.apply(this, tsFileProcessor, sequence);
}
- }
- private void tryToUpdateInsertLastCache(InsertRowNode node) {
- if (!CommonDescriptor.getInstance().getConfig().isLastCacheEnable()
- ||
(config.getDataRegionConsensusProtocolClass().equals(ConsensusFactory.IOT_CONSENSUS)
- && node.isSyncFromLeaderWhenUsingIoTConsensus())) {
+ if (CommonDescriptor.getInstance().getConfig().isLastCacheEnable()) {
+ if
((config.getDataRegionConsensusProtocolClass().equals(ConsensusFactory.IOT_CONSENSUS)
+ && insertRowNode.isSyncFromLeaderWhenUsingIoTConsensus())) {
+ return;
+ }
// disable updating last cache on follower
- return;
+ long startTime = System.nanoTime();
+ tryToUpdateInsertRowLastCache(insertRowNode);
+
PERFORMANCE_OVERVIEW_METRICS.recordScheduleUpdateLastCacheCost(System.nanoTime()
- startTime);
}
+ }
+
+ private void tryToUpdateInsertRowLastCache(InsertRowNode node) {
long latestFlushedTime =
lastFlushTimeMap.getGlobalFlushedTime(node.getDevicePath().getFullPath());
String[] measurements = node.getMeasurements();
@@ -1164,6 +1170,84 @@ public class DataRegion implements IDataRegionForQuery {
latestFlushedTime);
}
+ private void insertToTsFileProcessors(
+ InsertRowsNode insertRowsNode, boolean[] areSequence, long[]
timePartitionIds) {
+ List<InsertRowNode> executedInsertRowNodeList = new ArrayList<>();
+ long[] costsForMetrics = new long[4];
+ for (int i = 0; i < areSequence.length; i++) {
+ InsertRowNode insertRowNode =
insertRowsNode.getInsertRowNodeList().get(i);
+ if (insertRowNode.allMeasurementFailed()) {
+ continue;
+ }
+ TsFileProcessor tsFileProcessor =
+ getOrCreateTsFileProcessor(timePartitionIds[i], areSequence[i]);
+ if (tsFileProcessor == null) {
+ continue;
+ }
+ 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]);
+ }
+ }
+
+
PERFORMANCE_OVERVIEW_METRICS.recordCreateMemtableBlockCost(costsForMetrics[0]);
+
PERFORMANCE_OVERVIEW_METRICS.recordScheduleMemoryBlockCost(costsForMetrics[1]);
+ PERFORMANCE_OVERVIEW_METRICS.recordScheduleWalCost(costsForMetrics[2]);
+
PERFORMANCE_OVERVIEW_METRICS.recordScheduleMemTableCost(costsForMetrics[3]);
+
+ if (CommonDescriptor.getInstance().getConfig().isLastCacheEnable()) {
+ if
((config.getDataRegionConsensusProtocolClass().equals(ConsensusFactory.IOT_CONSENSUS)
+ && insertRowsNode.isSyncFromLeaderWhenUsingIoTConsensus())) {
+ return;
+ }
+ // disable updating last cache on follower
+ long startTime = System.nanoTime();
+ tryToUpdateInsertRowsLastCache(executedInsertRowNodeList);
+
PERFORMANCE_OVERVIEW_METRICS.recordScheduleUpdateLastCacheCost(System.nanoTime()
- startTime);
+ }
+ }
+
+ private void tryToUpdateInsertRowsLastCache(List<InsertRowNode> nodeList) {
+ DataNodeSchemaCache.getInstance().takeReadLock();
+ try {
+ for (InsertRowNode node : nodeList) {
+ long latestFlushedTime =
+
lastFlushTimeMap.getGlobalFlushedTime(node.getDevicePath().getFullPath());
+ String[] measurements = node.getMeasurements();
+ MeasurementSchema[] measurementSchemas = node.getMeasurementSchemas();
+ String[] rawMeasurements = new String[measurements.length];
+ for (int i = 0; i < measurements.length; i++) {
+ if (measurementSchemas[i] != null) {
+ // get raw measurement rather than alias
+ rawMeasurements[i] = measurementSchemas[i].getMeasurementId();
+ } else {
+ rawMeasurements[i] = measurements[i];
+ }
+ }
+ DataNodeSchemaCache.getInstance()
+ .updateLastCacheWithoutLock(
+ getDatabaseName(),
+ node.getDevicePath(),
+ rawMeasurements,
+ node.getMeasurementSchemas(),
+ node.isAligned(),
+ node::composeTimeValuePair,
+ index -> node.getValues()[index] != null,
+ true,
+ latestFlushedTime);
+ }
+ } finally {
+ DataNodeSchemaCache.getInstance().releaseReadLock();
+ }
+ }
+
/**
* WAL module uses this method to flush memTable
*
@@ -3029,23 +3113,62 @@ public class DataRegion implements IDataRegionForQuery {
}
}
- /**
- * insert batch of rows belongs to multiple devices
- *
- * @param insertRowsNode batch of rows belongs to multiple devices
- */
- public void insert(InsertRowsNode insertRowsNode) throws
BatchProcessException {
- for (int i = 0; i < insertRowsNode.getInsertRowNodeList().size(); i++) {
- InsertRowNode insertRowNode =
insertRowsNode.getInsertRowNodeList().get(i);
- try {
- insert(insertRowNode);
- } catch (WriteProcessException e) {
- insertRowsNode.getResults().put(i,
RpcUtils.getStatus(e.getErrorCode(), e.getMessage()));
- }
+ public void insert(InsertRowsNode insertRowsNode)
+ throws BatchProcessException, WriteProcessRejectException {
+ if (enableMemControl) {
+ StorageEngine.blockInsertionIfReject(null);
}
+ long startTime = System.nanoTime();
+ writeLock("InsertRows");
+ PERFORMANCE_OVERVIEW_METRICS.recordScheduleLockCost(System.nanoTime() -
startTime);
+ try {
+ if (deleted) {
+ return;
+ }
+ boolean[] areSequence = new
boolean[insertRowsNode.getInsertRowNodeList().size()];
+ long[] timePartitionIds = new
long[insertRowsNode.getInsertRowNodeList().size()];
+ for (int i = 0; i < insertRowsNode.getInsertRowNodeList().size(); i++) {
+ InsertRowNode insertRowNode =
insertRowsNode.getInsertRowNodeList().get(i);
+ if (!isAlive(insertRowNode.getTime())) {
+ insertRowsNode
+ .getResults()
+ .put(
+ i,
+ RpcUtils.getStatus(
+ TSStatusCode.OUT_OF_TTL.getStatusCode(),
+ String.format(
+ "Insertion time [%s] is less than ttl time bound
[%s]",
+
DateTimeUtils.convertLongToDate(insertRowNode.getTime()),
+ DateTimeUtils.convertLongToDate(
+ CommonDateTimeUtils.currentTime() - dataTTL))));
+ continue;
+ }
+ // init map
+ timePartitionIds[i] =
TimePartitionUtils.getTimePartitionId(insertRowNode.getTime());
- if (!insertRowsNode.getResults().isEmpty()) {
- throw new BatchProcessException("Partial failed inserting rows");
+ if
(!lastFlushTimeMap.checkAndCreateFlushedTimePartition(timePartitionIds[i])) {
+ TimePartitionManager.getInstance()
+ .registerTimePartitionInfo(
+ new TimePartitionInfo(
+ new DataRegionId(Integer.parseInt(dataRegionId)),
+ timePartitionIds[i],
+ true,
+ Long.MAX_VALUE,
+ 0,
+
tsFileManager.isLatestTimePartition(timePartitionIds[i])));
+ }
+ areSequence[i] =
+ config.isEnableSeparateData()
+ && insertRowNode.getTime()
+ > lastFlushTimeMap.getFlushedTime(
+ timePartitionIds[i],
insertRowNode.getDevicePath().getFullPath());
+ }
+ insertToTsFileProcessors(insertRowsNode, areSequence, timePartitionIds);
+ if (!insertRowsNode.getResults().isEmpty()) {
+ throw new BatchProcessException("Partial failed inserting rows");
+ }
+ } finally {
+ writeUnlock();
}
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java
index 33d160c3c37..837b4dd3491 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java
@@ -238,12 +238,14 @@ public class TsFileProcessor {
*
* @param insertRowNode physical plan of insertion
*/
- public void insert(InsertRowNode insertRowNode) throws WriteProcessException
{
+ public void insert(InsertRowNode insertRowNode, long[] costsForMetrics)
+ throws WriteProcessException {
if (workMemTable == null) {
long startTime = System.nanoTime();
createNewWorkingMemTable();
-
PERFORMANCE_OVERVIEW_METRICS.recordCreateMemtableBlockCost(System.nanoTime() -
startTime);
+ // recordCreateMemtableBlockCost
+ costsForMetrics[0] += System.nanoTime() - startTime;
WritingMetrics.getInstance()
.recordActiveMemTableCount(dataRegionInfo.getDataRegion().getDataRegionId(), 1);
}
@@ -262,7 +264,8 @@ public class TsFileProcessor {
insertRowNode.getDevicePath().getFullPath(),
insertRowNode.getMeasurements(),
insertRowNode.getDataTypes(), insertRowNode.getValues());
}
-
PERFORMANCE_OVERVIEW_METRICS.recordScheduleMemoryBlockCost(System.nanoTime() -
startTime);
+ // recordScheduleMemoryBlockCost
+ costsForMetrics[1] += System.nanoTime() - startTime;
}
long startTime = System.nanoTime();
@@ -282,7 +285,8 @@ public class TsFileProcessor {
storageGroupName, tsFileResource.getTsFile().getAbsolutePath()),
e);
} finally {
- PERFORMANCE_OVERVIEW_METRICS.recordScheduleWalCost(System.nanoTime() -
startTime);
+ // recordScheduleWalCost
+ costsForMetrics[2] += System.nanoTime() - startTime;
}
startTime = System.nanoTime();
@@ -315,8 +319,8 @@ public class TsFileProcessor {
}
tsFileResource.updateProgressIndex(insertRowNode.getProgressIndex());
-
- PERFORMANCE_OVERVIEW_METRICS.recordScheduleMemTableCost(System.nanoTime()
- startTime);
+ // recordScheduleMemTableCost
+ costsForMetrics[3] += System.nanoTime() - startTime;
}
private void createNewWorkingMemTable() throws WriteProcessException {
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessorTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessorTest.java
index 1315ef4cc84..1abb9cec674 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessorTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessorTest.java
@@ -129,7 +129,7 @@ public class TsFileProcessorTest {
for (int i = 1; i <= 100; i++) {
TSRecord record = new TSRecord(i, deviceId);
record.addTuple(DataPoint.getDataPoint(dataType, measurementId,
String.valueOf(i)));
- processor.insert(buildInsertRowNodeByTSRecord(record));
+ processor.insert(buildInsertRowNodeByTSRecord(record), new long[4]);
}
// query data in memory
@@ -188,7 +188,7 @@ public class TsFileProcessorTest {
for (int i = 1; i <= 100; i++) {
TSRecord record = new TSRecord(i, deviceId);
record.addTuple(DataPoint.getDataPoint(dataType, measurementId,
String.valueOf(i)));
- processor.insert(buildInsertRowNodeByTSRecord(record));
+ processor.insert(buildInsertRowNodeByTSRecord(record), new long[4]);
}
// query data in memory
@@ -273,7 +273,7 @@ public class TsFileProcessorTest {
for (int i = 1; i <= 10; i++) {
TSRecord record = new TSRecord(i, deviceId);
record.addTuple(DataPoint.getDataPoint(dataType, measurementId,
String.valueOf(i)));
- processor.insert(buildInsertRowNodeByTSRecord(record));
+ processor.insert(buildInsertRowNodeByTSRecord(record), new long[4]);
}
processor.asyncFlush();
}
@@ -315,7 +315,7 @@ public class TsFileProcessorTest {
for (int i = 1; i <= 100; i++) {
TSRecord record = new TSRecord(i, deviceId);
record.addTuple(DataPoint.getDataPoint(dataType, measurementId,
String.valueOf(i)));
- processor.insert(buildInsertRowNodeByTSRecord(record));
+ processor.insert(buildInsertRowNodeByTSRecord(record), new long[4]);
}
Assert.assertEquals(1598424, memTable.getTVListsRamCost());
Assert.assertEquals(90100, memTable.getTotalPointsNum());
@@ -362,14 +362,14 @@ public class TsFileProcessorTest {
for (int i = 1; i <= 100; i++) {
TSRecord record = new TSRecord(i, deviceId);
record.addTuple(DataPoint.getDataPoint(dataType, measurementId,
String.valueOf(i)));
- processor.insert(buildInsertRowNodeByTSRecord(record));
+ processor.insert(buildInsertRowNodeByTSRecord(record), new long[4]);
}
Assert.assertEquals(6387232, memTable.getTVListsRamCost());
// Test records
for (int i = 1; i <= 100; i++) {
TSRecord record = new TSRecord(i, deviceId);
record.addTuple(DataPoint.getDataPoint(dataType, "s1",
String.valueOf(i)));
- processor.insert(buildInsertRowNodeByTSRecord(record));
+ processor.insert(buildInsertRowNodeByTSRecord(record), new long[4]);
}
Assert.assertEquals(6388848, memTable.getTVListsRamCost());
Assert.assertEquals(240200, memTable.getTotalPointsNum());
@@ -405,7 +405,7 @@ public class TsFileProcessorTest {
for (int i = 1; i <= 100; i++) {
TSRecord record = new TSRecord(i, deviceId);
record.addTuple(DataPoint.getDataPoint(dataType, measurementId,
String.valueOf(i)));
- processor.insert(buildInsertRowNodeByTSRecord(record));
+ processor.insert(buildInsertRowNodeByTSRecord(record), new long[4]);
}
Assert.assertEquals(3193616, memTable.getTVListsRamCost());
Assert.assertEquals(90100, memTable.getTotalPointsNum());
@@ -442,7 +442,7 @@ public class TsFileProcessorTest {
for (int i = 1; i <= 100; i++) {
TSRecord record = new TSRecord(i, deviceId);
record.addTuple(DataPoint.getDataPoint(dataType, measurementId,
String.valueOf(i)));
- processor.insert(buildInsertRowNodeByTSRecord(record));
+ processor.insert(buildInsertRowNodeByTSRecord(record), new long[4]);
}
// query data in memory