This is an automated email from the ASF dual-hosted git repository.
rong 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 ecacaa1be88 Load: Use IDeviceID in ChunkData (#13887)
ecacaa1be88 is described below
commit ecacaa1be88054ec62bd6c0ab56b4aede7536ac8
Author: Zikun Ma <[email protected]>
AuthorDate: Thu Oct 24 11:33:13 2024 +0800
Load: Use IDeviceID in ChunkData (#13887)
---
.../plan/scheduler/load/LoadTsFileScheduler.java | 6 +---
.../db/storageengine/load/LoadTsFileManager.java | 6 ++--
.../load/splitter/AlignedChunkData.java | 33 ++++++++++++++--------
.../splitter/BatchedAlignedValueChunkData.java | 3 +-
.../db/storageengine/load/splitter/ChunkData.java | 5 ++--
.../load/splitter/NonAlignedChunkData.java | 28 +++++++++++-------
.../load/splitter/TsFileSplitter.java | 9 ++----
.../BatchedCompactionWithTsFileSplitterTest.java | 3 +-
8 files changed, 52 insertions(+), 41 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java
index cfa4965fe5d..5fc847b2b9b 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java
@@ -591,11 +591,7 @@ public class LoadTsFileScheduler implements IScheduler {
List<TRegionReplicaSet> replicaSets =
scheduler.partitionFetcher.queryDataPartition(
nonDirectionalChunkData.stream()
- .map(
- data ->
- new Pair<>(
-
IDeviceID.Factory.DEFAULT_FACTORY.create(data.getDevice()),
- data.getTimePartitionSlot()))
+ .map(data -> new Pair<>(data.getDevice(),
data.getTimePartitionSlot()))
.collect(Collectors.toList()),
scheduler.queryContext.getSession().getUserName());
IntStream.range(0, nonDirectionalChunkData.size())
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/LoadTsFileManager.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/LoadTsFileManager.java
index 75c82a9a1d3..775dad1df46 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/LoadTsFileManager.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/LoadTsFileManager.java
@@ -361,7 +361,7 @@ public class LoadTsFileManager {
private final File taskDir;
private Map<DataPartitionInfo, TsFileIOWriter> dataPartition2Writer;
- private Map<DataPartitionInfo, String> dataPartition2LastDevice;
+ private Map<DataPartitionInfo, IDeviceID> dataPartition2LastDevice;
private Map<DataPartitionInfo, ModificationFile>
dataPartition2ModificationFile;
private boolean isClosed;
@@ -408,11 +408,11 @@ public class LoadTsFileManager {
dataPartition2Writer.put(partitionInfo, writer);
}
TsFileIOWriter writer = dataPartition2Writer.get(partitionInfo);
- if
(!chunkData.getDevice().equals(dataPartition2LastDevice.getOrDefault(partitionInfo,
""))) {
+ if (!Objects.equals(chunkData.getDevice(),
dataPartition2LastDevice.get(partitionInfo))) {
if (dataPartition2LastDevice.containsKey(partitionInfo)) {
writer.endChunkGroup();
}
-
writer.startChunkGroup(IDeviceID.Factory.DEFAULT_FACTORY.create(chunkData.getDevice()));
+ writer.startChunkGroup(chunkData.getDevice());
dataPartition2LastDevice.put(partitionInfo, chunkData.getDevice());
}
chunkData.writeToFileWriter(writer);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/AlignedChunkData.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/AlignedChunkData.java
index f20a34cb3ca..585421e3afe 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/AlignedChunkData.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/AlignedChunkData.java
@@ -22,17 +22,18 @@ package org.apache.iotdb.db.storageengine.load.splitter;
import org.apache.iotdb.common.rpc.thrift.TTimePartitionSlot;
import org.apache.iotdb.commons.utils.TimePartitionUtils;
-import org.apache.tsfile.common.conf.TSFileConfig;
import org.apache.tsfile.enums.TSDataType;
import org.apache.tsfile.exception.write.PageException;
import org.apache.tsfile.file.header.ChunkHeader;
import org.apache.tsfile.file.header.PageHeader;
import org.apache.tsfile.file.metadata.IChunkMetadata;
+import org.apache.tsfile.file.metadata.IDeviceID;
+import org.apache.tsfile.file.metadata.PlainDeviceID;
+import org.apache.tsfile.file.metadata.StringArrayDeviceID;
import org.apache.tsfile.file.metadata.statistics.Statistics;
import org.apache.tsfile.read.common.Chunk;
import org.apache.tsfile.utils.Binary;
import org.apache.tsfile.utils.PublicBAOS;
-import org.apache.tsfile.utils.ReadWriteForEncodingUtils;
import org.apache.tsfile.utils.ReadWriteIOUtils;
import org.apache.tsfile.utils.TsPrimitiveType;
import org.apache.tsfile.write.UnSupportedDataTypeException;
@@ -41,6 +42,8 @@ import org.apache.tsfile.write.schema.IMeasurementSchema;
import org.apache.tsfile.write.schema.MeasurementSchema;
import org.apache.tsfile.write.writer.TsFileIOWriter;
+import javax.annotation.Nonnull;
+
import java.io.DataOutputStream;
import java.io.IOException;
import java.io.InputStream;
@@ -60,7 +63,7 @@ public class AlignedChunkData implements ChunkData {
protected static final Binary DEFAULT_BINARY = null;
protected final TTimePartitionSlot timePartitionSlot;
- protected final String device;
+ protected final IDeviceID device;
protected List<ChunkHeader> chunkHeaderList;
protected final PublicBAOS byteStream;
@@ -75,7 +78,7 @@ public class AlignedChunkData implements ChunkData {
protected List<Chunk> chunkList;
public AlignedChunkData(
- final String device,
+ @Nonnull final IDeviceID device,
final ChunkHeader chunkHeader,
final TTimePartitionSlot timePartitionSlot) {
this(device, timePartitionSlot);
@@ -84,14 +87,15 @@ public class AlignedChunkData implements ChunkData {
addAttrDataSize();
}
- protected AlignedChunkData(AlignedChunkData alignedChunkData) {
+ protected AlignedChunkData(final AlignedChunkData alignedChunkData) {
this(alignedChunkData.device, alignedChunkData.timePartitionSlot);
this.satisfiedLengthQueue = new
LinkedList<>(alignedChunkData.satisfiedLengthQueue);
this.needDecodeChunk = alignedChunkData.needDecodeChunk;
addAttrDataSize();
}
- protected AlignedChunkData(String device, TTimePartitionSlot
timePartitionSlot) {
+ protected AlignedChunkData(
+ @Nonnull final IDeviceID device, final TTimePartitionSlot
timePartitionSlot) {
this.dataSize = 0;
this.device = device;
this.chunkHeaderList = new ArrayList<>();
@@ -106,9 +110,7 @@ public class AlignedChunkData implements ChunkData {
private void addAttrDataSize() { // Should be init before serialize,
corresponding serializeAttr
dataSize += 2 * Byte.BYTES; // isModification and isAligned
dataSize += Long.BYTES; // timePartitionSlot
- final int deviceLength =
device.getBytes(TSFileConfig.STRING_CHARSET).length;
- dataSize += ReadWriteForEncodingUtils.varIntSize(deviceLength);
- dataSize += deviceLength; // device
+ dataSize += device.serializedSize(); // device
dataSize += Integer.BYTES; // chunkHeaderListSize
if (!chunkHeaderList.isEmpty()) {
dataSize += chunkHeaderList.get(0).getSerializedSize(); //
timeChunkHeader
@@ -116,7 +118,7 @@ public class AlignedChunkData implements ChunkData {
}
@Override
- public String getDevice() {
+ public IDeviceID getDevice() {
return device;
}
@@ -171,7 +173,10 @@ public class AlignedChunkData implements ChunkData {
private void serializeAttr(final DataOutputStream stream) throws IOException
{
ReadWriteIOUtils.write(timePartitionSlot.getStartTime(), stream);
- ReadWriteIOUtils.write(device, stream);
+
+ ReadWriteIOUtils.write(device instanceof StringArrayDeviceID, stream);
+ device.serialize(stream);
+
ReadWriteIOUtils.write(dataSize, stream);
ReadWriteIOUtils.write(needDecodeChunk, stream);
ReadWriteIOUtils.write(chunkHeaderList.size(), stream);
@@ -421,7 +426,11 @@ public class AlignedChunkData implements ChunkData {
throws IOException, PageException {
final TTimePartitionSlot timePartitionSlot =
TimePartitionUtils.getTimePartitionSlot(ReadWriteIOUtils.readLong(stream));
- final String device = ReadWriteIOUtils.readString(stream);
+ final boolean isStringArrayDeviceID = ReadWriteIOUtils.readBool(stream);
+ final IDeviceID device =
+ isStringArrayDeviceID
+ ? StringArrayDeviceID.deserialize(stream)
+ : PlainDeviceID.deserialize(stream).convertToStringArrayDeviceId();
final long dataSize = ReadWriteIOUtils.readLong(stream);
final boolean needDecodeChunk = ReadWriteIOUtils.readBool(stream);
final int chunkHeaderListSize = ReadWriteIOUtils.readInt(stream);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/BatchedAlignedValueChunkData.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/BatchedAlignedValueChunkData.java
index 2cac3414696..3bf50a9a296 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/BatchedAlignedValueChunkData.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/BatchedAlignedValueChunkData.java
@@ -26,6 +26,7 @@ import org.apache.tsfile.enums.TSDataType;
import org.apache.tsfile.exception.write.PageException;
import org.apache.tsfile.file.header.ChunkHeader;
import org.apache.tsfile.file.header.PageHeader;
+import org.apache.tsfile.file.metadata.IDeviceID;
import org.apache.tsfile.file.metadata.statistics.Statistics;
import org.apache.tsfile.read.common.Chunk;
import org.apache.tsfile.utils.Binary;
@@ -59,7 +60,7 @@ public class BatchedAlignedValueChunkData extends
AlignedChunkData {
}
// Used for deserialize
- public BatchedAlignedValueChunkData(String device, TTimePartitionSlot
timePartitionSlot) {
+ public BatchedAlignedValueChunkData(IDeviceID device, TTimePartitionSlot
timePartitionSlot) {
super(device, timePartitionSlot);
valueChunkWriters = new ArrayList<>();
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/ChunkData.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/ChunkData.java
index f64121527aa..1e96fac3868 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/ChunkData.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/ChunkData.java
@@ -25,6 +25,7 @@ import org.apache.tsfile.exception.write.PageException;
import org.apache.tsfile.file.header.ChunkHeader;
import org.apache.tsfile.file.header.PageHeader;
import org.apache.tsfile.file.metadata.IChunkMetadata;
+import org.apache.tsfile.file.metadata.IDeviceID;
import org.apache.tsfile.utils.ReadWriteIOUtils;
import org.apache.tsfile.write.writer.TsFileIOWriter;
@@ -33,7 +34,7 @@ import java.io.InputStream;
import java.nio.ByteBuffer;
public interface ChunkData extends TsFileData {
- String getDevice();
+ IDeviceID getDevice();
TTimePartitionSlot getTimePartitionSlot();
@@ -63,7 +64,7 @@ public interface ChunkData extends TsFileData {
static ChunkData createChunkData(
boolean isAligned,
- String device,
+ IDeviceID device,
ChunkHeader chunkHeader,
TTimePartitionSlot timePartitionSlot) {
return isAligned
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/NonAlignedChunkData.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/NonAlignedChunkData.java
index 230e7c127c5..f2d17ab95fc 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/NonAlignedChunkData.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/NonAlignedChunkData.java
@@ -22,22 +22,25 @@ package org.apache.iotdb.db.storageengine.load.splitter;
import org.apache.iotdb.common.rpc.thrift.TTimePartitionSlot;
import org.apache.iotdb.commons.utils.TimePartitionUtils;
-import org.apache.tsfile.common.conf.TSFileConfig;
import org.apache.tsfile.exception.write.PageException;
import org.apache.tsfile.file.header.ChunkHeader;
import org.apache.tsfile.file.header.PageHeader;
import org.apache.tsfile.file.metadata.IChunkMetadata;
+import org.apache.tsfile.file.metadata.IDeviceID;
+import org.apache.tsfile.file.metadata.PlainDeviceID;
+import org.apache.tsfile.file.metadata.StringArrayDeviceID;
import org.apache.tsfile.file.metadata.statistics.Statistics;
import org.apache.tsfile.read.common.Chunk;
import org.apache.tsfile.utils.Binary;
import org.apache.tsfile.utils.PublicBAOS;
-import org.apache.tsfile.utils.ReadWriteForEncodingUtils;
import org.apache.tsfile.utils.ReadWriteIOUtils;
import org.apache.tsfile.write.UnSupportedDataTypeException;
import org.apache.tsfile.write.chunk.ChunkWriterImpl;
import org.apache.tsfile.write.schema.MeasurementSchema;
import org.apache.tsfile.write.writer.TsFileIOWriter;
+import javax.annotation.Nonnull;
+
import java.io.DataOutputStream;
import java.io.IOException;
import java.io.InputStream;
@@ -47,7 +50,7 @@ import java.nio.ByteBuffer;
public class NonAlignedChunkData implements ChunkData {
private final TTimePartitionSlot timePartitionSlot;
- private final String device;
+ private final IDeviceID device;
private final ChunkHeader chunkHeader;
private final PublicBAOS byteStream;
@@ -60,7 +63,7 @@ public class NonAlignedChunkData implements ChunkData {
private Chunk chunk;
public NonAlignedChunkData(
- final String device,
+ @Nonnull final IDeviceID device,
final ChunkHeader chunkHeader,
final TTimePartitionSlot timePartitionSlot) {
this.dataSize = 0;
@@ -78,14 +81,12 @@ public class NonAlignedChunkData implements ChunkData {
private void addAttrDataSize() { // should be init before serialize,
corresponding serializeAttr
dataSize += 2 * Byte.BYTES; // isModification and isAligned
dataSize += Long.BYTES; // timePartitionSlot
- final int deviceLength =
device.getBytes(TSFileConfig.STRING_CHARSET).length;
- dataSize += ReadWriteForEncodingUtils.varIntSize(deviceLength);
- dataSize += deviceLength; // device
+ dataSize += device.serializedSize(); // device
dataSize += chunkHeader.getSerializedSize(); // timeChunkHeader
}
@Override
- public String getDevice() {
+ public IDeviceID getDevice() {
return device;
}
@@ -129,7 +130,10 @@ public class NonAlignedChunkData implements ChunkData {
private void serializeAttr(final DataOutputStream stream) throws IOException
{
ReadWriteIOUtils.write(timePartitionSlot.getStartTime(), stream);
- ReadWriteIOUtils.write(device, stream);
+
+ ReadWriteIOUtils.write(device instanceof StringArrayDeviceID, stream);
+ device.serialize(stream);
+
ReadWriteIOUtils.write(dataSize, stream);
ReadWriteIOUtils.write(needDecodeChunk, stream);
chunkHeader.serializeTo(stream); // chunk header already serialize chunk
type
@@ -275,7 +279,11 @@ public class NonAlignedChunkData implements ChunkData {
throws IOException, PageException {
final TTimePartitionSlot timePartitionSlot =
TimePartitionUtils.getTimePartitionSlot(ReadWriteIOUtils.readLong(stream));
- final String device = ReadWriteIOUtils.readString(stream);
+ final boolean isStringArrayDeviceID = ReadWriteIOUtils.readBool(stream);
+ final IDeviceID device =
+ isStringArrayDeviceID
+ ? StringArrayDeviceID.deserialize(stream)
+ : PlainDeviceID.deserialize(stream).convertToStringArrayDeviceId();
final long dataSize = ReadWriteIOUtils.readLong(stream);
final boolean needDecodeChunk = ReadWriteIOUtils.readBool(stream);
final byte chunkType = ReadWriteIOUtils.readByte(stream);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/TsFileSplitter.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/TsFileSplitter.java
index f2536f94568..e657ac0fa50 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/TsFileSplitter.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/splitter/TsFileSplitter.java
@@ -173,7 +173,7 @@ public class TsFileSplitter {
TTimePartitionSlot timePartitionSlot =
TimePartitionUtils.getTimePartitionSlot(chunkMetadata.getStartTime());
ChunkData chunkData =
- ChunkData.createChunkData(isAligned, curDevice.toString(), header,
timePartitionSlot);
+ ChunkData.createChunkData(isAligned, curDevice, header,
timePartitionSlot);
if (!needDecodeChunk(chunkMetadata)) {
chunkData.setNotDecode();
@@ -230,8 +230,7 @@ public class TsFileSplitter {
consumeChunkData(measurementId, chunkOffset, chunkData);
}
timePartitionSlot = pageTimePartitionSlot;
- chunkData =
- ChunkData.createChunkData(isAligned, curDevice.toString(),
header, timePartitionSlot);
+ chunkData = ChunkData.createChunkData(isAligned, curDevice, header,
timePartitionSlot);
}
if (isAligned) {
pageIndex2ChunkData
@@ -267,9 +266,7 @@ public class TsFileSplitter {
satisfiedLength = 0;
endTime =
timePartitionSlot.getStartTime() +
TimePartitionUtils.getTimePartitionInterval();
- chunkData =
- ChunkData.createChunkData(
- isAligned, curDevice.toString(), header,
timePartitionSlot);
+ chunkData = ChunkData.createChunkData(isAligned, curDevice,
header, timePartitionSlot);
}
satisfiedLength += 1;
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/BatchedCompactionWithTsFileSplitterTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/BatchedCompactionWithTsFileSplitterTest.java
index 6fca32c1277..9daee881371 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/BatchedCompactionWithTsFileSplitterTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/BatchedCompactionWithTsFileSplitterTest.java
@@ -278,8 +278,7 @@ public class BatchedCompactionWithTsFileSplitterTest
extends AbstractCompactionT
}
});
try {
- IDeviceID deviceID =
-
IDeviceID.Factory.DEFAULT_FACTORY.create(alignedChunkData.getDevice());
+ final IDeviceID deviceID = alignedChunkData.getDevice();
if (!deviceID.equals(writer.currentDevice)) {
if (writer.currentDevice != null) {
writer.endChunkGroup();