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 72aeed2591f Pipe: Optimized insert node cache hit possibility & Pipe: 
Set thread name for pipe receiver (#15263)
72aeed2591f is described below

commit 72aeed2591f23de2f3a567b43e12d02e5ee74ca4
Author: Steve Yurong Su <[email protected]>
AuthorDate: Wed Apr 2 20:00:20 2025 +0800

    Pipe: Optimized insert node cache hit possibility & Pipe: Set thread name 
for pipe receiver (#15263)
    
    * Pipe: Optimized insert node cache hit possibility
    
    * Pipe: Set thread name for pipe receiver
---
 .../db/storageengine/dataregion/wal/utils/WALEntryHandler.java | 10 ++++++++--
 .../storageengine/dataregion/wal/utils/WALEntryPosition.java   |  4 ++--
 .../storageengine/dataregion/wal/utils/WALInsertNodeCache.java |  6 ++++++
 .../apache/iotdb/commons/pipe/receiver/IoTDBFileReceiver.java  |  4 ++++
 4 files changed, 20 insertions(+), 4 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/utils/WALEntryHandler.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/utils/WALEntryHandler.java
index 333842e38ef..f5d7406f5a6 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/utils/WALEntryHandler.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/utils/WALEntryHandler.java
@@ -26,6 +26,7 @@ import 
org.apache.iotdb.db.storageengine.dataregion.wal.exception.MemTablePinExc
 import 
org.apache.iotdb.db.storageengine.dataregion.wal.exception.WALPipeException;
 import org.apache.iotdb.db.storageengine.dataregion.wal.node.WALNode;
 
+import org.apache.tsfile.utils.Pair;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -90,8 +91,13 @@ public class WALEntryHandler {
   public InsertNode getInsertNodeViaCacheIfPossible() {
     try {
       final WALEntryValue finalValue = value;
-      return finalValue instanceof InsertNode ? (InsertNode) finalValue : null;
-    } catch (Exception e) {
+      if (finalValue instanceof InsertNode) {
+        return (InsertNode) finalValue;
+      }
+      final Pair<ByteBuffer, InsertNode> byteBufferInsertNodePair =
+          walEntryPosition.getByteBufferOrInsertNodeIfPossible();
+      return byteBufferInsertNodePair == null ? null : 
byteBufferInsertNodePair.getRight();
+    } catch (final Exception e) {
       logger.warn("Fail to get insert node via cache. {}", this, e);
       throw e;
     }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/utils/WALEntryPosition.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/utils/WALEntryPosition.java
index 438e6a897e2..1370745cea2 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/utils/WALEntryPosition.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/utils/WALEntryPosition.java
@@ -64,8 +64,8 @@ public class WALEntryPosition {
    * Try to read the wal entry directly from the cache. No need to check if 
the wal entry is ready
    * for read.
    */
-  public Pair<ByteBuffer, InsertNode> 
readByteBufferOrInsertNodeViaCacheDirectly() {
-    return cache.getByteBufferOrInsertNode(this);
+  public Pair<ByteBuffer, InsertNode> getByteBufferOrInsertNodeIfPossible() {
+    return cache.getByteBufferOrInsertNodeIfPossible(this);
   }
 
   /**
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/utils/WALInsertNodeCache.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/utils/WALInsertNodeCache.java
index 21571c1ed98..0823c7e7b6e 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/utils/WALInsertNodeCache.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/utils/WALInsertNodeCache.java
@@ -229,6 +229,12 @@ public class WALInsertNodeCache {
     return pair;
   }
 
+  public Pair<ByteBuffer, InsertNode> getByteBufferOrInsertNodeIfPossible(
+      final WALEntryPosition position) {
+    hasPipeRunning = true;
+    return lruCache.getIfPresent(position);
+  }
+
   public void cacheInsertNodeIfNeeded(
       final WALEntryPosition walEntryPosition, final InsertNode insertNode) {
     // reduce memory usage
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/receiver/IoTDBFileReceiver.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/receiver/IoTDBFileReceiver.java
index a733fcde0a6..a8c65a93808 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/receiver/IoTDBFileReceiver.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/receiver/IoTDBFileReceiver.java
@@ -115,6 +115,10 @@ public abstract class IoTDBFileReceiver implements 
IoTDBReceiver {
     }
 
     receiverId.set(RECEIVER_ID_GENERATOR.incrementAndGet());
+    Thread.currentThread()
+        .setName(
+            String.format(
+                "Pipe-Receiver-%s-%s:%s", receiverId.get(), getSenderHost(), 
getSenderPort()));
 
     // Clear the original receiver file dir if exists
     if (receiverFileDirWithIdSuffix.get() != null) {

Reply via email to