This is an automated email from the ASF dual-hosted git repository.

haonan pushed a commit to branch object_type
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/object_type by this push:
     new d851d0fd3c9 impl wal consensus
d851d0fd3c9 is described below

commit d851d0fd3c9773d3a76cb3b5e32fd63c2cc6aeba
Author: HTHou <[email protected]>
AuthorDate: Mon Jul 7 19:09:03 2025 +0800

    impl wal consensus
---
 .../plan/planner/plan/node/write/FileNode.java     | 44 +++++++++++++++++++---
 .../node/write/RelationalInsertTabletNode.java     |  7 +++-
 .../dataregion/wal/buffer/WALBuffer.java           |  3 ++
 .../dataregion/wal/buffer/WALInfoEntry.java        |  2 +-
 .../storageengine/dataregion/wal/node/WALNode.java | 11 +++++-
 5 files changed, 58 insertions(+), 9 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/FileNode.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/FileNode.java
index 55fbf337d39..cdd393d579f 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/FileNode.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/FileNode.java
@@ -21,6 +21,7 @@ package 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write;
 
 import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
 import org.apache.iotdb.commons.consensus.index.ProgressIndex;
+import org.apache.iotdb.commons.exception.ObjectFileNotExist;
 import org.apache.iotdb.db.queryengine.plan.analyze.IAnalysis;
 import org.apache.iotdb.db.queryengine.plan.planner.plan.node.PlanNode;
 import org.apache.iotdb.db.queryengine.plan.planner.plan.node.PlanNodeId;
@@ -29,14 +30,18 @@ import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.WritePlanNode;
 import 
org.apache.iotdb.db.storageengine.dataregion.wal.buffer.IWALByteBufferView;
 import org.apache.iotdb.db.storageengine.dataregion.wal.buffer.WALEntryValue;
 import org.apache.iotdb.db.storageengine.dataregion.wal.utils.WALWriteUtils;
+import org.apache.iotdb.db.storageengine.rescon.disk.TierManager;
 
 import org.apache.tsfile.utils.ReadWriteIOUtils;
 
 import java.io.DataInputStream;
 import java.io.DataOutputStream;
+import java.io.File;
 import java.io.IOException;
+import java.io.RandomAccessFile;
 import java.nio.ByteBuffer;
 import java.util.List;
+import java.util.Optional;
 
 // TODO:[OBJECT] WAL serde
 public class FileNode extends SearchNode implements WALEntryValue {
@@ -56,7 +61,7 @@ public class FileNode extends SearchNode implements 
WALEntryValue {
     this.content = content;
   }
 
-  public FileNode(boolean isEOF, long offset, String filePath) {
+  public FileNode(boolean isEOF, long offset, byte[] content, String filePath) 
{
     super(new PlanNodeId(""));
     this.isEOF = isEOF;
     this.offset = offset;
@@ -89,6 +94,7 @@ public class FileNode extends SearchNode implements 
WALEntryValue {
     buffer.putLong(searchIndex);
     buffer.put((byte) (isEOF ? 1 : 0));
     buffer.putLong(offset);
+    buffer.putInt(content.length);
     WALWriteUtils.write(filePath, buffer);
   }
 
@@ -98,6 +104,7 @@ public class FileNode extends SearchNode implements 
WALEntryValue {
         + Long.BYTES
         + Byte.BYTES
         + Long.BYTES
+        + Integer.BYTES
         + ReadWriteIOUtils.sizeToWrite(filePath);
   }
 
@@ -105,10 +112,22 @@ public class FileNode extends SearchNode implements 
WALEntryValue {
     long searchIndex = stream.readLong();
     boolean isEOF = stream.readByte() == 1;
     long offset = stream.readLong();
+    int contentLength = stream.readInt();
     String filePath = ReadWriteIOUtils.readString(stream);
-
-    FileNode fileNode = new FileNode(isEOF, offset, filePath);
+    Optional<File> objectFile = 
TierManager.getInstance().getAbsoluteObjectFilePath(filePath);
+    byte[] contents = new byte[contentLength];
+    if (objectFile.isPresent()) {
+      try (RandomAccessFile raf = new RandomAccessFile(filePath, "r")) {
+        raf.seek(offset);
+        raf.read(contents);
+      }
+    } else {
+      throw new ObjectFileNotExist(filePath);
+    }
+
+    FileNode fileNode = new FileNode(isEOF, offset, contents, filePath);
     fileNode.setSearchIndex(searchIndex);
+
     return fileNode;
   }
 
@@ -116,9 +135,22 @@ public class FileNode extends SearchNode implements 
WALEntryValue {
     long searchIndex = buffer.getLong();
     boolean isEOF = buffer.get() == 1;
     long offset = buffer.getLong();
+    int contentLength = buffer.getInt();
     String filePath = ReadWriteIOUtils.readString(buffer);
-
-    FileNode fileNode = new FileNode(isEOF, offset, filePath);
+    Optional<File> objectFile = 
TierManager.getInstance().getAbsoluteObjectFilePath(filePath);
+    byte[] contents = new byte[contentLength];
+    if (objectFile.isPresent()) {
+      try (RandomAccessFile raf = new RandomAccessFile(filePath, "r")) {
+        raf.seek(offset);
+        raf.read(contents);
+      } catch (IOException e) {
+        throw new RuntimeException(e);
+      }
+    } else {
+      throw new ObjectFileNotExist(filePath);
+    }
+
+    FileNode fileNode = new FileNode(isEOF, offset, contents, filePath);
     fileNode.setSearchIndex(searchIndex);
     return fileNode;
   }
@@ -182,6 +214,6 @@ public class FileNode extends SearchNode implements 
WALEntryValue {
 
   @Override
   public long getMemorySize() {
-    return super.getMemorySize();
+    return content.length;
   }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertTabletNode.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertTabletNode.java
index 68fb6c4fbe9..03a67dc9516 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertTabletNode.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertTabletNode.java
@@ -396,6 +396,9 @@ public class RelationalInsertTabletNode extends 
InsertTabletNode {
   }
 
   public void handleObjectTypeValue() {
+    if (isGeneratedByRemoteConsensusLeader) {
+      return;
+    }
     List<List<FileNode>> fileNodesList = new ArrayList<>();
     for (int i = 0; i < dataTypes.length; i++) {
       if (dataTypes[i] == TSDataType.OBJECT) {
@@ -416,6 +419,8 @@ public class RelationalInsertTabletNode extends 
InsertTabletNode {
         fileNodesList.add(fileNodes);
       }
     }
-    this.fileNodesList = fileNodesList;
+    if (!fileNodesList.isEmpty()) {
+      this.fileNodesList = fileNodesList;
+    }
   }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALBuffer.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALBuffer.java
index 16be3f8ad08..f0a95a97705 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALBuffer.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALBuffer.java
@@ -26,6 +26,7 @@ import org.apache.iotdb.commons.utils.TestOnly;
 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.DeleteDataNode;
+import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.FileNode;
 import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertNode;
 import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.RelationalDeleteDataNode;
 import org.apache.iotdb.db.service.metrics.WritingMetrics;
@@ -332,6 +333,8 @@ public class WALBuffer extends AbstractWALBuffer {
           searchIndex = ((DeleteDataNode) 
walEntry.getValue()).getSearchIndex();
         } else if (walEntry.getType() == 
WALEntryType.RELATIONAL_DELETE_DATA_NODE) {
           searchIndex = ((RelationalDeleteDataNode) 
walEntry.getValue()).getSearchIndex();
+        } else if (walEntry.getType() == WALEntryType.OBJECT_FILE_NODE) {
+          searchIndex = ((FileNode) walEntry.getValue()).getSearchIndex();
         } else {
           searchIndex = ((InsertNode) walEntry.getValue()).getSearchIndex();
         }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALInfoEntry.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALInfoEntry.java
index edfdf0411a2..57157a0092b 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALInfoEntry.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALInfoEntry.java
@@ -169,7 +169,7 @@ public class WALInfoEntry extends WALEntry {
       case MEMORY_TABLE_CHECKPOINT:
         return RamUsageEstimator.sizeOfObject(value);
       case OBJECT_FILE_NODE:
-        return ((FileNode) value).getMemorySize();
+        return ((FileNode) value).serializedSize();
       default:
         throw new RuntimeException("Unsupported wal entry type " + type);
     }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java
index 56551e898aa..dafc852fd44 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java
@@ -65,6 +65,8 @@ import org.apache.tsfile.utils.TsFileUtils;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.io.ByteArrayInputStream;
+import java.io.DataInputStream;
 import java.io.File;
 import java.io.FileNotFoundException;
 import java.io.IOException;
@@ -775,7 +777,14 @@ public class WALNode implements IWALNode {
               } else if (currentWalEntryIndex < nextSearchIndex) {
                 // WAL entry is outdated, do nothing, continue to see next WAL 
entry
               } else if (currentWalEntryIndex == nextSearchIndex) {
-                tmpNodes.get().add(new IoTConsensusRequest(buffer));
+                if (type == WALEntryType.OBJECT_FILE_NODE) {
+                  WALEntry walEntry =
+                      WALEntry.deserialize(
+                          new DataInputStream(new 
ByteArrayInputStream(buffer.array())));
+                  tmpNodes.get().add((FileNode) walEntry.getValue());
+                } else {
+                  tmpNodes.get().add(new IoTConsensusRequest(buffer));
+                }
               } else {
                 // currentWalEntryIndex > targetIndex
                 // WAL entry of targetIndex has been fully collected, put them 
into insertNodes

Reply via email to