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