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();

Reply via email to