Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -21,8 +21,10 @@

import org.apache.iotdb.db.conf.IoTDBConfig;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowNode;
import org.apache.iotdb.db.storageengine.dataregion.wal.buffer.IWALByteBufferView;
import org.apache.iotdb.db.storageengine.dataregion.wal.utils.WALWriteUtils;
import org.apache.iotdb.db.utils.MemUtils;
import org.apache.iotdb.db.utils.datastructure.AlignedTVList;
import org.apache.iotdb.db.utils.datastructure.BatchEncodeInfo;
import org.apache.iotdb.db.utils.datastructure.MemPointIterator;
Expand All @@ -46,6 +48,7 @@
import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
Expand All @@ -58,6 +61,15 @@
import static org.apache.iotdb.db.storageengine.rescon.memory.PrimitiveArrayManager.ARRAY_SIZE;
import static org.apache.iotdb.db.utils.ModificationUtils.isPointDeleted;

interface AlignedRowsMemCostEstimator {

void addRow(InsertRowNode row);

long getTVListMemoryCost();

long getTextDataMemoryCost();
}

public class AlignedWritableMemChunk extends AbstractWritableMemChunk {

private final Map<String, Integer> measurementIndexMap;
Expand Down Expand Up @@ -104,6 +116,130 @@ public boolean containsMeasurement(String measurementId) {
return measurementIndexMap.containsKey(measurementId);
}

public long alignedWriteRowMemCost(
String[] measurements, TSDataType[] incomingDataTypes, Object[] rowValues, int rowOffset) {
return list.alignedWriteRowMemCost(
mapMeasurementsToTVListColumns(measurements), incomingDataTypes, rowValues, rowOffset);
}

public long alignedWriteArrayMemCost(
String[] measurements, TSDataType[] incomingDataTypes, Object[] columns, int start, int end) {
return list.alignedWriteArrayMemCost(
mapMeasurementsToTVListColumns(measurements), incomingDataTypes, columns, start, end);
}

public long alignedWriteRowsMemCost(List<InsertRowNode> rows) {
AlignedTVList workingList = list;
synchronized (workingList) {
AlignedRowsMemCostEstimator estimator =
new ExistingAlignedRowsMemCostEstimator(workingList, false);
for (InsertRowNode row : rows) {
estimator.addRow(row);
}
return estimator.getTVListMemoryCost();
}
}

// TsFileProcessor serializes writes for a DataRegion, so the incremental estimator can retain a
// stable view of this TVList while rows from multiple devices are processed in input order.
AlignedRowsMemCostEstimator newAlignedWriteRowsMemCostEstimator() {
return new ExistingAlignedRowsMemCostEstimator(list, true);
}

private final class ExistingAlignedRowsMemCostEstimator implements AlignedRowsMemCostEstimator {
private final AlignedTVList.AlignedWriteRowsMemCostEstimator tvListEstimator;
private final Map<String, Integer> newMeasurementIndexMap = new HashMap<>();
private final boolean trackTextData;

private int nextColumnIndex = dataTypes.size();
private String[] previousMeasurements;
private int[] previousColumnIndexes = new int[0];
private long textDataMemoryCost;

private ExistingAlignedRowsMemCostEstimator(AlignedTVList workingList, boolean trackTextData) {
tvListEstimator = workingList.newAlignedWriteRowsMemCostEstimator();
this.trackTextData = trackTextData;
}

@Override
public void addRow(InsertRowNode row) {
tvListEstimator.startRow();
String[] measurements = row.getMeasurements();
TSDataType[] incomingDataTypes = row.getDataTypes();
Object[] rowValues = row.getValues();
int[] tvListColumnIndexes = mapMeasurements(measurements);
int columnCount =
Math.min(measurements.length, Math.min(incomingDataTypes.length, rowValues.length));
for (int column = 0; column < columnCount; column++) {
String measurement = measurements[column];
TSDataType dataType = incomingDataTypes[column];
Object value = rowValues[column];
tvListEstimator.addValue(tvListColumnIndexes[column], dataType, value);
if (trackTextData
&& measurement != null
&& dataType != null
&& dataType.isBinary()
&& value != null) {
textDataMemoryCost += MemUtils.getBinarySize((Binary) value);
}
}
}

private int[] mapMeasurements(String[] measurements) {
if (previousMeasurements != null
&& (previousMeasurements == measurements
|| Arrays.equals(previousMeasurements, measurements))) {
previousMeasurements = measurements;
return previousColumnIndexes;
}
if (previousColumnIndexes.length < measurements.length) {
previousColumnIndexes = new int[measurements.length];
}
Arrays.fill(previousColumnIndexes, 0, measurements.length, -1);
for (int i = 0; i < measurements.length; i++) {
String measurement = measurements[i];
if (measurement == null) {
continue;
}
Integer columnIndex = measurementIndexMap.get(measurement);
if (columnIndex == null) {
columnIndex = newMeasurementIndexMap.get(measurement);
if (columnIndex == null) {
columnIndex = nextColumnIndex++;
newMeasurementIndexMap.put(measurement, columnIndex);
}
}
previousColumnIndexes[i] = columnIndex;
}
previousMeasurements = measurements;
return previousColumnIndexes;
}

@Override
public long getTVListMemoryCost() {
return tvListEstimator.getMemoryCost();
}

@Override
public long getTextDataMemoryCost() {
return textDataMemoryCost;
}
}

private int[] mapMeasurementsToTVListColumns(String[] measurements) {
int[] columnIndexes = new int[measurements.length];
Arrays.fill(columnIndexes, -1);
for (int i = 0; i < measurements.length; i++) {
if (measurements[i] != null) {
Integer columnIndex = measurementIndexMap.get(measurements[i]);
if (columnIndex != null) {
columnIndexes[i] = columnIndex;
}
}
}
return columnIndexes;
}

@Override
public void putLong(long t, long v) {
throw new UnSupportedDataTypeException(UNSUPPORTED_TYPE + TSDataType.VECTOR);
Expand Down
Loading
Loading